fix(recycle): fail-safe entity recycling + OSS batch-delete diff bug (#86/#3232)

Two blob-cleanup defects:

- OSS driver diffed `files` (the whole request) instead of `group` (the
  current 1000-key batch) when computing undeleted keys — every batch
  beyond the first was pre-marked failed, so deleted blobs kept their DB
  rows and the sweep re-tried them forever.
- A blob that can never be deleted (illegal over-long path, corrupted
  key) blocked the routine sweep indefinitely: every run re-attempted
  it, errored the task, and repeated. EntityProps gains a JSON
  `recycle_fail_count`; per-entity failures increment it (whole-chunk
  failures are treated as infra outages and never count), and past
  RecycleFailThreshold (5) the row is force-removed with a loud warning
  rather than stalling cleanup forever.

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 fb96eaebaa
commit 04583175bb

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

@ -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 {

@ -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

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

@ -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
}
// 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)
}
}
return entity.ID(), force
})
// 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)

Loading…
Cancel
Save