Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 0 additions & 1 deletion drivers/alias/driver.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
}
Expand Down
4 changes: 2 additions & 2 deletions internal/model/args.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"`
Expand All @@ -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,
Expand Down
25 changes: 25 additions & 0 deletions internal/model/args_test.go
Original file line number Diff line number Diff line change
@@ -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")
}
}
18 changes: 11 additions & 7 deletions internal/op/archive.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 (
Expand All @@ -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
}
}
Expand All @@ -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)
}
Expand All @@ -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
}
}
Expand All @@ -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 {
Expand Down
14 changes: 8 additions & 6 deletions internal/op/fs.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
}
}
Expand All @@ -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)
}
Expand All @@ -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
}
}
Expand Down
34 changes: 34 additions & 0 deletions internal/op/link_lifecycle.go
Original file line number Diff line number Diff line change
@@ -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
}
143 changes: 143 additions & 0 deletions internal/op/link_lifecycle_test.go
Original file line number Diff line number Diff line change
@@ -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())
}
})
}
Loading