diff --git a/inventory/entity_props_test.go b/inventory/entity_props_test.go new file mode 100644 index 00000000..c67c01ed --- /dev/null +++ b/inventory/entity_props_test.go @@ -0,0 +1,52 @@ +package inventory + +import ( + "context" + "testing" + + "github.com/cloudreve/Cloudreve/v4/ent/enttest" + "github.com/cloudreve/Cloudreve/v4/inventory/types" + "github.com/stretchr/testify/require" +) + +func TestUpdateEntityProps(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() + + policy := client.StoragePolicy.Create(). + SetName("local"). + SetType(string(types.PolicyTypeLocal)). + SetBucketName("bucket"). + SetAccessKey("ak"). + SetSecretKey("sk"). + SetMaxSize(0). + SetDirNameRule("d"). + SetFileNameRule("f"). + SaveX(ctx) + + e := client.Entity.Create(). + SetType(int(types.EntityTypeVersion)). + SetSource("uploads/1/test.bin"). + SetSize(100). + SetStoragePolicyEntities(policy.ID). + SaveX(ctx) + + c := NewFileClient(client, "sqlite3", nil) + + e.Props = &types.EntityProps{RecycleFailCount: 3} + require.NoError(t, c.UpdateEntityProps(ctx, e)) + + got := client.Entity.GetX(ctx, e.ID) + require.NotNil(t, got.Props) + require.Equal(t, 3, got.Props.RecycleFailCount) + + // Other props survive a partial update. + e.Props.RecycleFailCount = 5 + e.Props.UnlinkOnly = true + require.NoError(t, c.UpdateEntityProps(ctx, e)) + + got = client.Entity.GetX(ctx, e.ID) + require.Equal(t, 5, got.Props.RecycleFailCount) + require.True(t, got.Props.UnlinkOnly) +} diff --git a/inventory/file.go b/inventory/file.go index 3b344bd9..cc019cb9 100644 --- a/inventory/file.go +++ b/inventory/file.go @@ -187,6 +187,8 @@ type FileClient interface { RemoveStaleEntities(ctx context.Context, file *ent.File) (StorageDiff, error) // RemoveEntitiesByID hard-delete entities by IDs. RemoveEntitiesByID(ctx context.Context, ids ...int) (map[int]int64, error) + // UpdateEntityProps persists mutated EntityProps back to the given entities. + UpdateEntityProps(ctx context.Context, entities ...*ent.Entity) error // CapEntities caps the number of entities of a given file. The oldest entities will be unlinked // if entity count exceed limit. CapEntities(ctx context.Context, file *ent.File, owner *ent.User, max int, entityType types.EntityType) (StorageDiff, error) @@ -458,6 +460,15 @@ func (f *fileClient) RemoveEntitiesByID(ctx context.Context, ids ...int) (map[in return storageReduced, nil } +func (f *fileClient) UpdateEntityProps(ctx context.Context, entities ...*ent.Entity) error { + for _, e := range entities { + if err := f.client.Entity.UpdateOne(e).SetProps(e.Props).Exec(ctx); err != nil { + return fmt.Errorf("failed to update props of entity %d: %w", e.ID, err) + } + } + return nil +} + func (f *fileClient) StaleEntities(ctx context.Context, ids ...int) ([]*ent.Entity, error) { res := make([]*ent.Entity, 0, len(ids)) if len(ids) > 0 { diff --git a/inventory/types/types.go b/inventory/types/types.go index febd32a6..acc6a6bb 100644 --- a/inventory/types/types.go +++ b/inventory/types/types.go @@ -179,6 +179,10 @@ type ( EntityProps struct { UnlinkOnly bool `json:"unlink_only,omitempty"` EncryptMetadata *EncryptMetadata `json:"encrypt_metadata,omitempty"` + // RecycleFailCount tracks consecutive driver-delete failures during + // entity recycling. Entities reaching the threshold are force-removed + // from the DB so one un-deletable blob cannot stall the sweep forever. + RecycleFailCount int `json:"recycle_fail_count,omitempty"` } Cipher string diff --git a/pkg/filemanager/driver/oss/oss.go b/pkg/filemanager/driver/oss/oss.go index 1b382533..c0191e37 100644 --- a/pkg/filemanager/driver/oss/oss.go +++ b/pkg/filemanager/driver/oss/oss.go @@ -351,7 +351,7 @@ func (handler *Driver) Delete(ctx context.Context, files ...string) ([]string, e // 统计未删除的文件 failed = append( failed, - util.SliceDifference(files, + util.SliceDifference(group, lo.Map(delRes.DeletedObjects, func(v oss.DeletedInfo, i int) string { return *v.Key }), diff --git a/pkg/filemanager/manager/recycle.go b/pkg/filemanager/manager/recycle.go index e01abfbd..eb91d21f 100644 --- a/pkg/filemanager/manager/recycle.go +++ b/pkg/filemanager/manager/recycle.go @@ -40,6 +40,10 @@ type ( } ) +// RecycleFailThreshold is the number of consecutive sweeps an entity may fail +// driver deletion before its DB row is force-removed (blob orphaned). +const RecycleFailThreshold = 5 + func init() { queue.RegisterResumableTaskFactory(queue.ExplicitEntityRecycleTaskType, NewExplicitEntityRecycleTaskFromModel) queue.RegisterResumableTaskFactory(queue.EntityRecycleRoutineTaskType, NewEntityRecycleRoutineTaskFromModel) @@ -226,14 +230,19 @@ func (m *manager) RecycleEntities(ctx context.Context, force bool, entityIDs ... }), func(entity fs.Entity, index int) string { return entity.Source() }) + var failedSrcs []string if len(toBeDeletedSrc) > 0 { - res, err := d.Delete(ctx, toBeDeletedSrc...) - if err != nil { - for _, src := range res { - ae.Add(strconv.Itoa(mapSrcToId[src]), err) + var deleteErr error + failedSrcs, deleteErr = d.Delete(ctx, toBeDeletedSrc...) + if deleteErr != nil { + for _, src := range failedSrcs { + ae.Add(strconv.Itoa(mapSrcToId[src]), deleteErr) } } } + // A failure covering every source signals an infrastructure outage + // (driver/bucket down) — per-entity badness only counts on partial failure. + partialFailure := len(failedSrcs) < len(toBeDeletedSrc) // Delete upload session if it's still valid for _, entity := range chunk { @@ -253,25 +262,53 @@ func (m *manager) RecycleEntities(ctx context.Context, force bool, entityIDs ... // Filtering out entities that are successfully deleted rawAe := ae.Raw() - successEntities := lo.FilterMap(chunk, func(entity fs.Entity, index int) (int, bool) { + successEntities := make([]int, 0, len(chunk)) + var dirtyProps []*ent.Entity + for _, entity := range chunk { entityIdStr := fmt.Sprintf("%d", entity.ID()) - _, ok := rawAe[entityIdStr] - if !ok { - // No error, deleted - return entity.ID(), true + if _, ok := rawAe[entityIdStr]; !ok { + successEntities = append(successEntities, entity.ID()) + continue } if force { ae.Remove(entityIdStr) + successEntities = append(successEntities, entity.ID()) + continue } - return entity.ID(), force - }) + + // Count consecutive per-entity delete failures; past the + // threshold the row is force-removed so one un-deletable blob + // (e.g. an illegal over-long path) cannot stall every sweep. + if partialFailure { + model := entity.Model() + if model.Props == nil { + model.Props = &types.EntityProps{} + } + model.Props.RecycleFailCount++ + if model.Props.RecycleFailCount >= RecycleFailThreshold { + m.l.Warning( + "Entity %d failed blob deletion %d times; removing DB row and orphaning blob %q", + entity.ID(), model.Props.RecycleFailCount, entity.Source()) + ae.Remove(entityIdStr) + successEntities = append(successEntities, entity.ID()) + continue + } + dirtyProps = append(dirtyProps, model) + } + } // Remove entities from DB fc, tx, ctx, err := inventory.WithTx(ctx, m.dep.FileClient()) if err != nil { return fmt.Errorf("failed to start transaction: %w", err) } + if len(dirtyProps) > 0 { + if err := fc.UpdateEntityProps(ctx, dirtyProps...); err != nil { + _ = inventory.Rollback(tx) + return fmt.Errorf("failed to persist recycle failure counts: %w", err) + } + } storageReduced, err := fc.RemoveEntitiesByID(ctx, successEntities...) if err != nil { _ = inventory.Rollback(tx)