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 <info@tdvorak.dev>

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 9bf2e14b5f
commit 63066db78a

@ -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())
}

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