From 63066db78adb7854bd02b83dd9804a8826bcaa54 Mon Sep 17 00:00:00 2001 From: Tomas Dvorak Date: Fri, 18 Sep 2026 22:57:28 +0200 Subject: [PATCH] fix(recycle): let force delete bypass broken storage drivers (#2909) RecycleEntities aborted the whole chunk when getEntityPolicyDriver failed, so blob rows bound to a deleted or misconfigured storage policy could never be removed - which in turn blocked deleting the policy. Under force=true the driver-init failure is now logged and skipped (the physical blobs are unreachable anyway), and the DB rows are removed. The non-force path is unchanged: it still reports per-entity errors and retains the rows. Authored By: TDvorak Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- pkg/filemanager/manager/recycle.go | 17 +++- pkg/filemanager/manager/recycle_force_test.go | 83 +++++++++++++++++++ 2 files changed, 96 insertions(+), 4 deletions(-) create mode 100644 pkg/filemanager/manager/recycle_force_test.go diff --git a/pkg/filemanager/manager/recycle.go b/pkg/filemanager/manager/recycle.go index eb91d21f..7929c05f 100644 --- a/pkg/filemanager/manager/recycle.go +++ b/pkg/filemanager/manager/recycle.go @@ -213,12 +213,19 @@ func (m *manager) RecycleEntities(ctx context.Context, force bool, entityIDs ... m.l.Info("Start to recycle batch #%d, %d entities", batch, len(chunk)) mapSrcToId := make(map[string]int, len(chunk)) _, d, err := m.getEntityPolicyDriver(ctx, chunk[0], nil) - if err != nil { + if err != nil && !force { for _, entity := range chunk { ae.Add(strconv.Itoa(entity.ID()), err) } continue } + if err != nil { + // Force delete: the policy/driver is broken (e.g. deleted or + // misconfigured), so physical blobs are unreachable anyway — + // drop the DB rows instead of deadlocking the policy. + m.l.Warning("Skipping driver init failure under force recycle: %s", err) + d = nil + } for _, entity := range chunk { mapSrcToId[entity.Source()] = entity.ID() @@ -231,7 +238,7 @@ func (m *manager) RecycleEntities(ctx context.Context, force bool, entityIDs ... return entity.Source() }) var failedSrcs []string - if len(toBeDeletedSrc) > 0 { + if len(toBeDeletedSrc) > 0 && d != nil { var deleteErr error failedSrcs, deleteErr = d.Delete(ctx, toBeDeletedSrc...) if deleteErr != nil { @@ -253,8 +260,10 @@ func (m *manager) RecycleEntities(ctx context.Context, force bool, entityIDs ... if session, ok := m.kv.Get(UploadSessionCachePrefix + sid.String()); ok { session := session.(fs.UploadSession) - if err := d.CancelToken(ctx, &session); err != nil { - m.l.Warning("Failed to cancel upload session for %q: %s, this is expected if it's remote policy.", session.Props.Uri.String(), err) + if d != nil { + if err := d.CancelToken(ctx, &session); err != nil { + m.l.Warning("Failed to cancel upload session for %q: %s, this is expected if it's remote policy.", session.Props.Uri.String(), err) + } } _ = m.kv.Delete(UploadSessionCachePrefix, sid.String()) } diff --git a/pkg/filemanager/manager/recycle_force_test.go b/pkg/filemanager/manager/recycle_force_test.go new file mode 100644 index 00000000..bf8f44db --- /dev/null +++ b/pkg/filemanager/manager/recycle_force_test.go @@ -0,0 +1,83 @@ +package manager + +import ( + "context" + "path/filepath" + "testing" + + "github.com/cloudreve/Cloudreve/v4/application/dependency" + "github.com/cloudreve/Cloudreve/v4/ent/enttest" + "github.com/cloudreve/Cloudreve/v4/pkg/boolset" + "github.com/cloudreve/Cloudreve/v4/pkg/cache" + "github.com/cloudreve/Cloudreve/v4/pkg/conf" + "github.com/cloudreve/Cloudreve/v4/pkg/hashid" + "github.com/cloudreve/Cloudreve/v4/pkg/logging" + "github.com/stretchr/testify/require" +) + +// TestRecycleEntitiesForceBypassesBrokenPolicy covers upstream #2909: when a +// storage policy is deleted or too broken to initialize a driver, its blob +// rows used to be undeletable even with force=true, which also deadlocked +// removal of the policy itself. +func TestRecycleEntitiesForceBypassesBrokenPolicy(t *testing.T) { + client := enttest.Open(t, "sqlite3", "file:"+t.Name()+"?mode=memory&cache=shared") + t.Cleanup(func() { require.NoError(t, client.Close()) }) + + logger := logging.NewConsoleLogger(logging.LevelError) + cfg, err := conf.NewIniConfigProvider(filepath.Join(t.TempDir(), "conf.ini"), logger) + require.NoError(t, err) + hasher, err := hashid.New("test-salt") + require.NoError(t, err) + + dep := dependency.NewDependency( + dependency.WithDbClient(client), + dependency.WithConfigProvider(cfg), + dependency.WithKV(cache.NewMemoStore("", logger)), + dependency.WithLogger(logger), + dependency.WithHashIDEncoder(hasher), + ) + ctx := context.WithValue(context.Background(), dependency.DepCtx{}, dep) + + group := client.Group.Create().SetName("g").SetPermissions(&boolset.BooleanSet{}).SaveX(ctx) + user := client.User.Create().SetEmail("u@example.com").SetNick("u").SetGroup(group).SaveX(ctx) + + // Unknown policy type - GetStorageDriver rejects it with + // ErrUnknownPolicyType, simulating a misconfigured policy. + brokenPolicy := client.StoragePolicy.Create(). + SetName("broken"). + SetType("bogus"). + SetAccessKey("x"). + SetSecretKey("x"). + SetMaxSize(1). + SetDirNameRule("x"). + SetFileNameRule("x"). + SaveX(ctx) + + newStaleEntity := func() int { + return client.Entity.Create(). + SetType(0). + SetSource("uploads/blob.bin"). + SetSize(123). + SetStoragePolicyEntities(brokenPolicy.ID). + SetReferenceCount(0). + SetCreatedBy(user.ID). + SaveX(ctx).ID + } + + fm := NewFileManager(dep, user) + + t.Run("without force the row is retained and error returned", func(t *testing.T) { + id := newStaleEntity() + err := fm.RecycleEntities(ctx, false, id) + require.Error(t, err) + require.NotNil(t, client.Entity.GetX(ctx, id)) + }) + + t.Run("with force the row is dropped despite driver failure", func(t *testing.T) { + id := newStaleEntity() + err := fm.RecycleEntities(ctx, true, id) + require.NoError(t, err) + _, getErr := client.Entity.Get(ctx, id) + require.Error(t, getErr) + }) +}