From 56064d198165c58048c45944f59d53b710d964c8 Mon Sep 17 00:00:00 2001 From: Nostalgia Date: Mon, 21 Sep 2026 16:54:12 +0800 Subject: [PATCH] fix(op): enforce link cache lifecycle policy (#3101) - Share one admitted lifecycle policy between regular and archive links. - Reject and release links that combine TTL caching with owned resources. - Keep wrapper clones independent from the source cache expiration. - Cover reuse, reference, invalidation, conflict, and clone behavior. Co-authored-by: nostalume Co-authored-by: Codex <267193182+codex@users.noreply.github.com> --- drivers/alias/driver.go | 1 - internal/model/args.go | 4 +- internal/model/args_test.go | 25 +++++ internal/op/archive.go | 18 ++-- internal/op/fs.go | 14 +-- internal/op/link_lifecycle.go | 34 +++++++ internal/op/link_lifecycle_test.go | 143 +++++++++++++++++++++++++++++ 7 files changed, 223 insertions(+), 16 deletions(-) create mode 100644 internal/model/args_test.go create mode 100644 internal/op/link_lifecycle.go create mode 100644 internal/op/link_lifecycle_test.go diff --git a/drivers/alias/driver.go b/drivers/alias/driver.go index 7008495ce..d69d6cf50 100644 --- a/drivers/alias/driver.go +++ b/drivers/alias/driver.go @@ -328,7 +328,6 @@ func (d *Alias) Link(ctx context.Context, file model.Obj, args model.LinkArgs) ( return nil, err } resultLink := link.Clone() // 复制一份,避免修改到原始link - resultLink.Expiration = nil if args.Redirect { return resultLink, nil } diff --git a/internal/model/args.go b/internal/model/args.go index 393a716fa..b051106cc 100644 --- a/internal/model/args.go +++ b/internal/model/args.go @@ -30,7 +30,7 @@ type Link struct { Header http.Header `json:"header"` // needed header (for url) RangeReader RangeReaderIF `json:"-"` // recommended way if can't use URL - Expiration *time.Duration // local cache expire Duration + Expiration *time.Duration // local cache expiration; not transferred by Clone //for accelerating request, use multi-thread downloading Concurrency int `json:"concurrency"` @@ -42,12 +42,12 @@ type Link struct { RequireReference bool `json:"-"` } +// Clone transfers ownership of l without inheriting its cache expiration. func (l *Link) Clone() *Link { return &Link{ URL: l.URL, Header: l.Header, RangeReader: l.RangeReader, - Expiration: l.Expiration, Concurrency: l.Concurrency, PartSize: l.PartSize, ContentLength: l.ContentLength, diff --git a/internal/model/args_test.go b/internal/model/args_test.go new file mode 100644 index 000000000..a82ae44b6 --- /dev/null +++ b/internal/model/args_test.go @@ -0,0 +1,25 @@ +package model + +import ( + "testing" + "time" +) + +func TestLinkCloneTransfersOwnershipWithoutCachePolicy(t *testing.T) { + ttl := time.Minute + source := &Link{URL: "https://example.test/file", Expiration: &ttl} + + clone := source.Clone() + if clone.URL != source.URL { + t.Fatal("clone did not preserve transport data") + } + if clone.Expiration != nil { + t.Fatal("clone inherited source cache policy") + } + if err := clone.Close(); err != nil { + t.Fatal(err) + } + if !source.Expired() { + t.Fatal("closing clone did not release its source") + } +} diff --git a/internal/op/archive.go b/internal/op/archive.go index d0a53a919..bb3de11a1 100644 --- a/internal/op/archive.go +++ b/internal/op/archive.go @@ -390,8 +390,9 @@ func ArchiveGet(ctx context.Context, storage driver.Driver, path string, args mo } type objWithLink struct { - link *model.Link - obj model.Obj + link *model.Link + obj model.Obj + policy linkCachePolicy } var ( @@ -405,7 +406,7 @@ func DriverExtract(ctx context.Context, storage driver.Driver, path string, args } key := stdpath.Join(Key(storage, path), args.InnerPath) if ol, ok := extractCache.Get(key); ok { - if ol.link.Expiration != nil || ol.link.SyncClosers.AcquireReference() || !ol.link.RequireReference { + if ol.acquire() { return ol.link, ol.obj, nil } } @@ -415,8 +416,8 @@ func DriverExtract(ctx context.Context, storage driver.Driver, path string, args if err != nil { return nil, errors.Wrapf(err, "failed extract archive") } - if ol.link.Expiration != nil { - extractCache.SetWithTTL(key, ol, *ol.link.Expiration) + if ol.policy.expiration != nil { + extractCache.SetWithTTL(key, ol, *ol.policy.expiration) } else { extractCache.SetWithExpirable(key, ol, &ol.link.SyncClosers) } @@ -428,7 +429,7 @@ func DriverExtract(ctx context.Context, storage driver.Driver, path string, args if err != nil { return nil, nil, err } - if ol.link.SyncClosers.AcquireReference() || !ol.link.RequireReference { + if ol.acquire() { return ol.link, ol.obj, nil } } @@ -450,7 +451,10 @@ func driverExtract(ctx context.Context, storage driver.Driver, path string, args return nil, errors.WithStack(errs.NotFile) } link, err := storageAr.Extract(ctx, archiveFile, args) - return &objWithLink{link: link, obj: extracted}, err + if err != nil { + return nil, err + } + return admitLink(link, extracted) } type streamWithParent struct { diff --git a/internal/op/fs.go b/internal/op/fs.go index f82a3ca8f..033d89b96 100644 --- a/internal/op/fs.go +++ b/internal/op/fs.go @@ -242,8 +242,7 @@ func Link(ctx context.Context, storage driver.Driver, path string, args model.Li } key := Key(storage, path) if ol, exists := Cache.linkCache.GetType(key, typeKey); exists { - if ol.link.Expiration != nil || - ol.link.SyncClosers.AcquireReference() || !ol.link.RequireReference { + if ol.acquire() { return ol.link, ol.obj, nil } } @@ -261,9 +260,12 @@ func Link(ctx context.Context, storage driver.Driver, path string, args model.Li if err != nil { return nil, errors.Wrapf(err, "failed get link") } - ol := &objWithLink{link: link, obj: file} - if link.Expiration != nil { - Cache.linkCache.SetTypeWithTTL(key, typeKey, ol, *link.Expiration) + ol, err := admitLink(link, file) + if err != nil { + return nil, err + } + if ol.policy.expiration != nil { + Cache.linkCache.SetTypeWithTTL(key, typeKey, ol, *ol.policy.expiration) } else { Cache.linkCache.SetTypeWithExpirable(key, typeKey, ol, &link.SyncClosers) } @@ -274,7 +276,7 @@ func Link(ctx context.Context, storage driver.Driver, path string, args model.Li if err != nil { return nil, nil, err } - if ol.link.SyncClosers.AcquireReference() || !ol.link.RequireReference { + if ol.acquire() { return ol.link, ol.obj, nil } } diff --git a/internal/op/link_lifecycle.go b/internal/op/link_lifecycle.go new file mode 100644 index 000000000..71491c97b --- /dev/null +++ b/internal/op/link_lifecycle.go @@ -0,0 +1,34 @@ +package op + +import ( + "errors" + "time" + + "github.com/OpenListTeam/OpenList/v4/internal/model" +) + +var errConflictingLinkLifecycle = errors.New("invalid link lifecycle: expiration cannot be combined with owned closers or RequireReference") + +type linkCachePolicy struct { + expiration *time.Duration + requireReference bool +} + +func admitLink(link *model.Link, obj model.Obj) (*objWithLink, error) { + if link.Expiration != nil && (link.RequireReference || link.SyncClosers.Length() > 0) { + return nil, errors.Join(errConflictingLinkLifecycle, link.Close()) + } + return &objWithLink{ + link: link, + obj: obj, + policy: linkCachePolicy{ + expiration: link.Expiration, + requireReference: link.RequireReference, + }, + }, nil +} + +func (ol *objWithLink) acquire() bool { + return ol.policy.expiration != nil || + ol.link.SyncClosers.AcquireReference() || !ol.policy.requireReference +} diff --git a/internal/op/link_lifecycle_test.go b/internal/op/link_lifecycle_test.go new file mode 100644 index 000000000..53bb0baf7 --- /dev/null +++ b/internal/op/link_lifecycle_test.go @@ -0,0 +1,143 @@ +package op + +import ( + "context" + "strings" + "sync/atomic" + "testing" + "time" + + "github.com/OpenListTeam/OpenList/v4/internal/driver" + "github.com/OpenListTeam/OpenList/v4/internal/model" + "github.com/OpenListTeam/OpenList/v4/pkg/singleflight" + "github.com/OpenListTeam/OpenList/v4/pkg/utils" +) + +type linkLifecycleDriver struct { + model.Storage + links func() *model.Link + calls atomic.Int32 +} + +func (d *linkLifecycleDriver) Config() driver.Config { return driver.Config{} } +func (d *linkLifecycleDriver) GetAddition() driver.Additional { return nil } +func (d *linkLifecycleDriver) Init(context.Context) error { return nil } +func (d *linkLifecycleDriver) Drop(context.Context) error { return nil } +func (d *linkLifecycleDriver) List(context.Context, model.Obj, model.ListArgs) ([]model.Obj, error) { + return nil, nil +} +func (d *linkLifecycleDriver) Get(context.Context, string) (model.Obj, error) { + return &model.Object{Name: "file", Path: "/file"}, nil +} +func (d *linkLifecycleDriver) Link(context.Context, model.Obj, model.LinkArgs) (*model.Link, error) { + d.calls.Add(1) + return d.links(), nil +} + +func resetLinkLifecycleState(t *testing.T) { + t.Helper() + oldCache := Cache + Cache, linkG = NewCacheManager(), singleflight.Group[*objWithLink]{} + t.Cleanup(func() { Cache, linkG = oldCache, singleflight.Group[*objWithLink]{} }) +} + +func acquireTestLink(t *testing.T, d *linkLifecycleDriver) *model.Link { + t.Helper() + link, _, err := Link(context.Background(), d, "/file", model.LinkArgs{}) + if err != nil { + t.Fatal(err) + } + return link +} + +func TestLinkLifecycleModes(t *testing.T) { + t.Run("TTL descriptor remains reusable after close", func(t *testing.T) { + resetLinkLifecycleState(t) + ttl := time.Minute + d := &linkLifecycleDriver{ + Storage: model.Storage{MountPath: "/ttl"}, + links: func() *model.Link { return &model.Link{URL: "https://example.test/file", Expiration: &ttl} }, + } + + first := acquireTestLink(t, d) + _ = first.Close() + second := acquireTestLink(t, d) + if second.URL != first.URL || d.calls.Load() != 1 { + t.Fatalf("TTL link was not reused: calls=%d", d.calls.Load()) + } + _ = second.Close() + }) + + t.Run("references keep shared resources alive until final close", func(t *testing.T) { + resetLinkLifecycleState(t) + var closes atomic.Int32 + d := &linkLifecycleDriver{ + Storage: model.Storage{MountPath: "/reference"}, + links: func() *model.Link { + return &model.Link{ + URL: "https://example.test/file", + SyncClosers: utils.NewSyncClosers(utils.CloseFunc(func() error { closes.Add(1); return nil })), + RequireReference: true, + } + }, + } + + first := acquireTestLink(t, d) + second := acquireTestLink(t, d) + _ = first.Close() + if closes.Load() != 0 { + t.Fatal("shared resource closed while another reference was active") + } + _ = second.Close() + if closes.Load() != 1 { + t.Fatalf("final close count = %d, want 1", closes.Load()) + } + third := acquireTestLink(t, d) + _ = third.Close() + if d.calls.Load() != 2 || closes.Load() != 2 { + t.Fatalf("stale link was not replaced: calls=%d closes=%d", d.calls.Load(), closes.Load()) + } + }) + + t.Run("close-invalidated link is reacquired", func(t *testing.T) { + resetLinkLifecycleState(t) + d := &linkLifecycleDriver{ + Storage: model.Storage{MountPath: "/close-invalidated"}, + links: func() *model.Link { + return &model.Link{SyncClosers: utils.NewSyncClosers(utils.CloseFunc(func() error { return nil }))} + }, + } + + first := acquireTestLink(t, d) + _ = first.Close() + second := acquireTestLink(t, d) + _ = second.Close() + if d.calls.Load() != 2 { + t.Fatalf("driver calls = %d, want 2", d.calls.Load()) + } + }) + + t.Run("TTL with owned resources is rejected and released", func(t *testing.T) { + resetLinkLifecycleState(t) + ttl := time.Minute + var closes atomic.Int32 + d := &linkLifecycleDriver{ + Storage: model.Storage{MountPath: "/conflict"}, + links: func() *model.Link { + return &model.Link{ + Expiration: &ttl, + SyncClosers: utils.NewSyncClosers(utils.CloseFunc(func() error { closes.Add(1); return nil })), + RequireReference: true, + } + }, + } + + _, _, err := Link(context.Background(), d, "/file", model.LinkArgs{}) + if err == nil || !strings.Contains(err.Error(), "expiration cannot be combined") { + t.Fatalf("unexpected error: %v", err) + } + if closes.Load() != 1 { + t.Fatalf("rejected link close count = %d, want 1", closes.Load()) + } + }) +}