mirror of
https://github.com/OpenListTeam/OpenList.git
synced 2026-10-10 04:53:09 +08:00
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 <nostalucent@gmail.com> Co-authored-by: Codex <267193182+codex@users.noreply.github.com>
This commit is contained in:
@@ -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
|
||||
}
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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")
|
||||
}
|
||||
}
|
||||
+11
-7
@@ -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 {
|
||||
|
||||
+8
-6
@@ -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
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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())
|
||||
}
|
||||
})
|
||||
}
|
||||
Reference in New Issue
Block a user