From 80487122ea5fb326fc4070b80023aad5de335d99 Mon Sep 17 00:00:00 2001 From: MadDogOwner Date: Thu, 16 Jul 2026 22:08:57 +0800 Subject: [PATCH] refactor(aws-sdk)!: migrate to transfermanager Signed-off-by: MadDogOwner --- drivers/123/driver.go | 12 ++++-------- drivers/halalcloud/driver.go | 16 +++++++++------- drivers/mediatrack/driver.go | 17 +++++++++-------- drivers/s3/driver.go | 16 ++++++++-------- drivers/teambition/util.go | 17 +++++++++-------- drivers/thunder/driver.go | 10 ++++------ drivers/thunder_browser/driver.go | 10 ++++------ drivers/thunderx/driver.go | 10 ++++------ go.mod | 1 - 9 files changed, 51 insertions(+), 58 deletions(-) diff --git a/drivers/123/driver.go b/drivers/123/driver.go index 02b0dde62..3fa0e33ea 100644 --- a/drivers/123/driver.go +++ b/drivers/123/driver.go @@ -21,7 +21,7 @@ import ( "github.com/aws/aws-sdk-go-v2/aws" "github.com/aws/aws-sdk-go-v2/credentials" "github.com/aws/aws-sdk-go-v2/service/s3" - "github.com/aws/aws-sdk-go-v2/feature/s3/manager" + "github.com/aws/aws-sdk-go-v2/feature/s3/transfermanager" "github.com/go-resty/resty/v2" log "github.com/sirupsen/logrus" ) @@ -219,19 +219,15 @@ func (d *Pan123) Put(ctx context.Context, dstDir model.Obj, file model.FileStrea o.BaseEndpoint = aws.String(resp.Data.EndPoint) o.UsePathStyle = true }) - uploader := manager.NewUploader(s3Client) - if file.GetSize() > int64(manager.MaxUploadParts)*manager.DefaultUploadPartSize { - uploader.PartSize = file.GetSize() / int64(manager.MaxUploadParts-1) - } - input := &s3.PutObjectInput{ + tmClient := transfermanager.New(s3Client) + _, err = tmClient.UploadObject(ctx, &transfermanager.UploadObjectInput{ Bucket: &resp.Data.Bucket, Key: &resp.Data.Key, Body: driver.NewLimitedUploadStream(ctx, &driver.ReaderUpdatingProgress{ Reader: file, UpdateProgress: up, }), - } - _, err = uploader.Upload(ctx, input) + }) if err != nil { return err } diff --git a/drivers/halalcloud/driver.go b/drivers/halalcloud/driver.go index b7c30c2b1..c2ab209d2 100644 --- a/drivers/halalcloud/driver.go +++ b/drivers/halalcloud/driver.go @@ -15,11 +15,12 @@ import ( "github.com/OpenListTeam/OpenList/v4/internal/model" "github.com/OpenListTeam/OpenList/v4/internal/op" "github.com/OpenListTeam/OpenList/v4/internal/stream" + "github.com/OpenListTeam/OpenList/v4/pkg/utils" "github.com/OpenListTeam/OpenList/v4/pkg/http_range" "github.com/aws/aws-sdk-go-v2/aws" "github.com/aws/aws-sdk-go-v2/credentials" "github.com/aws/aws-sdk-go-v2/service/s3" - "github.com/aws/aws-sdk-go-v2/feature/s3/manager" + "github.com/aws/aws-sdk-go-v2/feature/s3/transfermanager" "github.com/city404/v6-public-rpc-proto/go/v6/common" pbPublicUser "github.com/city404/v6-public-rpc-proto/go/v6/user" pubUserFile "github.com/city404/v6-public-rpc-proto/go/v6/userfile" @@ -381,14 +382,15 @@ func (d *HalalCloud) put(ctx context.Context, dstDir model.Obj, fileStream model o.BaseEndpoint = aws.String(result.Endpoint) o.UsePathStyle = true }) - uploader := manager.NewUploader(s3Client, func(u *manager.Uploader) { - u.Concurrency = d.uploadThread + tmClient := transfermanager.New(s3Client, func(o *transfermanager.Options) { + if fileStream.GetSize() > int64(8*utils.MB)*10000 { + o.PartSizeBytes = fileStream.GetSize() / 9999 + } + o.Concurrency = d.uploadThread }) - if fileStream.GetSize() > int64(manager.MaxUploadParts)*manager.DefaultUploadPartSize { - uploader.PartSize = fileStream.GetSize() / int64(manager.MaxUploadParts-1) - } + reader := driver.NewLimitedUploadStream(ctx, fileStream) - _, err = uploader.Upload(ctx, &s3.PutObjectInput{ + _, err = tmClient.UploadObject(ctx, &transfermanager.UploadObjectInput{ Bucket: aws.String(result.Bucket), Key: aws.String(result.Key), Body: io.TeeReader(reader, driver.NewProgress(fileStream.GetSize(), up)), diff --git a/drivers/mediatrack/driver.go b/drivers/mediatrack/driver.go index 29e2497ae..017b670f2 100644 --- a/drivers/mediatrack/driver.go +++ b/drivers/mediatrack/driver.go @@ -17,7 +17,7 @@ import ( "github.com/aws/aws-sdk-go-v2/aws" "github.com/aws/aws-sdk-go-v2/credentials" "github.com/aws/aws-sdk-go-v2/service/s3" - "github.com/aws/aws-sdk-go-v2/feature/s3/manager" + "github.com/aws/aws-sdk-go-v2/feature/s3/transfermanager" "github.com/go-resty/resty/v2" "github.com/google/uuid" log "github.com/sirupsen/logrus" @@ -184,11 +184,13 @@ func (d *MediaTrack) Put(ctx context.Context, dstDir model.Obj, file model.FileS if err != nil { return err } - uploader := manager.NewUploader(s3Client) - if file.GetSize() > int64(manager.MaxUploadParts)*manager.DefaultUploadPartSize { - uploader.PartSize = file.GetSize() / int64(manager.MaxUploadParts-1) - } - input := &s3.PutObjectInput{ + tmClient := transfermanager.New(s3Client, func(o *transfermanager.Options) { + if file.GetSize() > int64(8*utils.MB)*10000 { + o.PartSizeBytes = file.GetSize() / 9999 + } + }) + + _, err = tmClient.UploadObject(ctx, &transfermanager.UploadObjectInput{ Bucket: &resp.Data.Bucket, Key: &resp.Data.Object, Body: driver.NewLimitedUploadStream(ctx, &driver.ReaderUpdatingProgress{ @@ -198,8 +200,7 @@ func (d *MediaTrack) Put(ctx context.Context, dstDir model.Obj, file model.FileS }, UpdateProgress: up, }), - } - _, err = uploader.Upload(ctx, input) + }) if err != nil { return err } diff --git a/drivers/s3/driver.go b/drivers/s3/driver.go index e490c3d24..383a905f5 100644 --- a/drivers/s3/driver.go +++ b/drivers/s3/driver.go @@ -20,7 +20,7 @@ import ( "github.com/aws/smithy-go" "github.com/aws/aws-sdk-go-v2/service/s3" "github.com/aws/aws-sdk-go-v2/service/s3/types" - "github.com/aws/aws-sdk-go-v2/feature/s3/manager" + "github.com/aws/aws-sdk-go-v2/feature/s3/transfermanager" "github.com/pkg/errors" log "github.com/sirupsen/logrus" ) @@ -212,14 +212,15 @@ func (d *S3) Remove(ctx context.Context, obj model.Obj) error { } func (d *S3) Put(ctx context.Context, dstDir model.Obj, s model.FileStreamer, up driver.UpdateProgress) error { - uploader := manager.NewUploader(d.client) - if s.GetSize() > int64(manager.MaxUploadParts)*manager.DefaultUploadPartSize { - uploader.PartSize = s.GetSize() / int64(manager.MaxUploadParts-1) - } key := getKey(stdpath.Join(dstDir.GetPath(), s.GetName()), false) contentType := s.GetMimetype() log.Debugln("key:", key) - input := &s3.PutObjectInput{ + tmClient := transfermanager.New(d.client, func(o *transfermanager.Options) { + if s.GetSize() > int64(8*utils.MB)*10000 { + o.PartSizeBytes = s.GetSize() / 9999 + } + }) + _, err := tmClient.UploadObject(ctx, &transfermanager.UploadObjectInput{ Bucket: &d.Bucket, Key: &key, Body: driver.NewLimitedUploadStream(ctx, &driver.ReaderUpdatingProgress{ @@ -227,8 +228,7 @@ func (d *S3) Put(ctx context.Context, dstDir model.Obj, s model.FileStreamer, up UpdateProgress: up, }), ContentType: &contentType, - } - _, err := uploader.Upload(ctx, input) + }) return err } diff --git a/drivers/teambition/util.go b/drivers/teambition/util.go index 9a1d32c67..3dd42b9e3 100644 --- a/drivers/teambition/util.go +++ b/drivers/teambition/util.go @@ -17,7 +17,7 @@ import ( "github.com/aws/aws-sdk-go-v2/aws" "github.com/aws/aws-sdk-go-v2/credentials" "github.com/aws/aws-sdk-go-v2/service/s3" - "github.com/aws/aws-sdk-go-v2/feature/s3/manager" + "github.com/aws/aws-sdk-go-v2/feature/s3/transfermanager" "github.com/go-resty/resty/v2" log "github.com/sirupsen/logrus" ) @@ -248,11 +248,13 @@ func (d *Teambition) newUpload(ctx context.Context, dstDir model.Obj, stream mod if err != nil { return err } - uploader := manager.NewUploader(s3Client) - if stream.GetSize() > int64(manager.MaxUploadParts)*manager.DefaultUploadPartSize { - uploader.PartSize = stream.GetSize() / int64(manager.MaxUploadParts-1) - } - input := &s3.PutObjectInput{ + tmClient := transfermanager.New(s3Client, func(o *transfermanager.Options) { + if stream.GetSize() > int64(8*utils.MB)*10000 { + o.PartSizeBytes = stream.GetSize() / 9999 + } + }) + + _, err = tmClient.UploadObject(ctx, &transfermanager.UploadObjectInput{ Bucket: &uploadToken.Upload.Bucket, Key: &uploadToken.Upload.Key, ContentDisposition: &uploadToken.Upload.ContentDisposition, @@ -261,8 +263,7 @@ func (d *Teambition) newUpload(ctx context.Context, dstDir model.Obj, stream mod Reader: stream, UpdateProgress: up, }), - } - _, err = uploader.Upload(ctx, input) + }) if err != nil { return err } diff --git a/drivers/thunder/driver.go b/drivers/thunder/driver.go index d3bdad7cd..c2267a601 100644 --- a/drivers/thunder/driver.go +++ b/drivers/thunder/driver.go @@ -19,7 +19,7 @@ import ( "github.com/aws/aws-sdk-go-v2/aws" "github.com/aws/aws-sdk-go-v2/credentials" "github.com/aws/aws-sdk-go-v2/service/s3" - "github.com/aws/aws-sdk-go-v2/feature/s3/manager" + "github.com/aws/aws-sdk-go-v2/feature/s3/transfermanager" "github.com/go-resty/resty/v2" ) @@ -416,11 +416,9 @@ func (xc *XunLeiCommon) Put(ctx context.Context, dstDir model.Obj, file model.Fi }, func(o *s3.Options) { o.BaseEndpoint = aws.String(param.Endpoint) }) - uploader := manager.NewUploader(s3Client) - if file.GetSize() > int64(manager.MaxUploadParts)*manager.DefaultUploadPartSize { - uploader.PartSize = file.GetSize() / int64(manager.MaxUploadParts-1) - } - _, err = uploader.Upload(ctx, &s3.PutObjectInput{ + tmClient := transfermanager.New(s3Client) + + _, err = tmClient.UploadObject(ctx, &transfermanager.UploadObjectInput{ Bucket: aws.String(param.Bucket), Key: aws.String(param.Key), Body: driver.NewLimitedUploadStream(ctx, &driver.ReaderUpdatingProgress{ diff --git a/drivers/thunder_browser/driver.go b/drivers/thunder_browser/driver.go index 51419ed4e..3e1543127 100644 --- a/drivers/thunder_browser/driver.go +++ b/drivers/thunder_browser/driver.go @@ -21,7 +21,7 @@ import ( "github.com/aws/aws-sdk-go-v2/aws" "github.com/aws/aws-sdk-go-v2/credentials" "github.com/aws/aws-sdk-go-v2/service/s3" - "github.com/aws/aws-sdk-go-v2/feature/s3/manager" + "github.com/aws/aws-sdk-go-v2/feature/s3/transfermanager" "github.com/go-resty/resty/v2" ) @@ -530,11 +530,9 @@ func (xc *XunLeiBrowserCommon) Put(ctx context.Context, dstDir model.Obj, stream }, func(o *s3.Options) { o.BaseEndpoint = aws.String(param.Endpoint) }) - uploader := manager.NewUploader(s3Client) - if stream.GetSize() > int64(manager.MaxUploadParts)*manager.DefaultUploadPartSize { - uploader.PartSize = stream.GetSize() / int64(manager.MaxUploadParts-1) - } - _, err = uploader.Upload(ctx, &s3.PutObjectInput{ + tmClient := transfermanager.New(s3Client) + + _, err = tmClient.UploadObject(ctx, &transfermanager.UploadObjectInput{ Bucket: aws.String(param.Bucket), Key: aws.String(param.Key), Expires: aws.Time(param.Expiration), diff --git a/drivers/thunderx/driver.go b/drivers/thunderx/driver.go index 0803a9c62..46f69073e 100644 --- a/drivers/thunderx/driver.go +++ b/drivers/thunderx/driver.go @@ -20,7 +20,7 @@ import ( "github.com/aws/aws-sdk-go-v2/aws" "github.com/aws/aws-sdk-go-v2/credentials" "github.com/aws/aws-sdk-go-v2/service/s3" - "github.com/aws/aws-sdk-go-v2/feature/s3/manager" + "github.com/aws/aws-sdk-go-v2/feature/s3/transfermanager" "github.com/go-resty/resty/v2" ) @@ -406,11 +406,9 @@ func (xc *XunLeiXCommon) Put(ctx context.Context, dstDir model.Obj, file model.F }, func(o *s3.Options) { o.BaseEndpoint = aws.String(param.Endpoint) }) - uploader := manager.NewUploader(s3Client) - if file.GetSize() > int64(manager.MaxUploadParts)*manager.DefaultUploadPartSize { - uploader.PartSize = file.GetSize() / int64(manager.MaxUploadParts-1) - } - _, err = uploader.Upload(ctx, &s3.PutObjectInput{ + tmClient := transfermanager.New(s3Client) + + _, err = tmClient.UploadObject(ctx, &transfermanager.UploadObjectInput{ Bucket: aws.String(param.Bucket), Key: aws.String(param.Key), Body: driver.NewLimitedUploadStream(ctx, &driver.ReaderUpdatingProgress{ diff --git a/go.mod b/go.mod index aea9466cb..48abfeab1 100644 --- a/go.mod +++ b/go.mod @@ -200,7 +200,6 @@ require ( github.com/andreburgaud/crypt2go v1.8.0 // indirect github.com/andybalholm/brotli v1.2.2 // indirect github.com/aws/aws-sdk-go-v2/credentials v1.19.29 - github.com/aws/aws-sdk-go-v2/feature/s3/manager v1.22.13 github.com/aws/aws-sdk-go-v2/service/s3 v1.105.1 github.com/aws/smithy-go v1.27.3 github.com/axgle/mahonia v0.0.0-20180208002826-3358181d7394