Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
24 commits
Select commit Hold shift + click to select a range
9f16512
fix(drivers/139): update path handling for family
xrgzs Jul 4, 2026
1c8e268
perf(drivers/139): add retry go for family upload
xrgzs Jul 4, 2026
031a9ae
feat(drivers/139): add FamilyCloudHost and GroupCloudHost handling
xrgzs Jul 4, 2026
fcb1904
perf(drivers/139): add retry go for personalnew upload
xrgzs Jul 4, 2026
51d7d59
fix(drivers/139): use new batchpartinfos for remaining parts
xrgzs Jul 4, 2026
deef171
fix(driver/139): do not use path in id
xrgzs Jul 6, 2026
b5d6c6b
fix(driver/139): refactor RootPath handling
xrgzs Jul 6, 2026
a7a193f
fix(driver/139): remove root path stripping in Init method
xrgzs Jul 6, 2026
0e7d4c1
chore(drivers/139): reduce upload retry attempts to 3
xrgzs Jul 6, 2026
550e9c8
fix(drivers/139): add rate limiting for MetaPersonalNew upload
xrgzs Jul 14, 2026
73c6fe3
fix(drivers/139): fix section reader leak and seek check in retry block
xrgzs Jul 14, 2026
b9b9f54
fix(drivers/139): move GetSectionReader outside retry loop
xrgzs Jul 16, 2026
57e1523
feat(drivers/139): add UseOldStreamUpload option
xrgzs Jul 16, 2026
bd99844
fix(drivers/139): add more param for group/family put
xrgzs Jul 16, 2026
e947e3f
fix(drivers/139): correct group params for new put
xrgzs Jul 16, 2026
ad6dfd6
fix(drivers/139): update personalPost to use getPersonalCloudHost
xrgzs Jul 17, 2026
cd91c64
chore(drivers/139): correct typo
xrgzs Jul 17, 2026
3dd3d71
fix(drivers/139): pass full partInfos slice for batched upload
xrgzs Jul 21, 2026
7c12fbb
fix(drivers/139): avoid double-counting progress on upload retries
xrgzs Jul 21, 2026
972d566
fix(drivers/139): gate route-host checks by cloud type
xrgzs Jul 21, 2026
fd79574
fix(drivers/139): add ProviderRoot back for groupgetfiles
xrgzs Jul 21, 2026
ca48b73
fix(drivers/139): set path to 0 when group old upload
xrgzs Jul 21, 2026
97c24ca
fix(drivers/139): improve error handling for family root path retrieval
xrgzs Jul 22, 2026
dd56d49
refactor(drivers/139): update condition checks for CloudHost initiali…
xrgzs Jul 22, 2026
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
219 changes: 153 additions & 66 deletions drivers/139/driver.go
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ import (
"github.com/OpenListTeam/OpenList/v4/pkg/cron"
"github.com/OpenListTeam/OpenList/v4/pkg/utils"
"github.com/OpenListTeam/OpenList/v4/pkg/utils/random"
"github.com/avast/retry-go"
log "github.com/sirupsen/logrus"
)

Expand All @@ -29,7 +30,9 @@ type Yun139 struct {
Account string
ref *Yun139
PersonalCloudHost string
RootPath string
FamilyCloudHost string
GroupCloudHost string
ProviderRoot string
}

func (d *Yun139) Config() driver.Config {
Expand Down Expand Up @@ -73,14 +76,26 @@ func (d *Yun139) Init(ctx context.Context) error {
return err
}
for _, policyItem := range resp.Data.RoutePolicyList {
if policyItem.ModName == "personal" {
switch policyItem.ModName {
case "personal":
d.PersonalCloudHost = policyItem.HttpsUrl
break
case "group":
d.GroupCloudHost = policyItem.HttpsUrl
case "family":
d.FamilyCloudHost = policyItem.HttpsUrl
}
}
if len(d.PersonalCloudHost) == 0 {
return fmt.Errorf("PersonalCloudHost is empty")
}
if d.isGroup() || d.isFamily() {
if len(d.GroupCloudHost) == 0 {
return fmt.Errorf("GroupCloudHost is empty")
}
if len(d.FamilyCloudHost) == 0 {
return fmt.Errorf("FamilyCloudHost is empty")
}
}

d.cron = cron.NewCron(time.Hour * 12)
d.cron.Do(func() {
Expand Down Expand Up @@ -108,14 +123,17 @@ func (d *Yun139) Init(ctx context.Context) error {
return err
}
case MetaFamily:
// Attempt to obtain data.path as the root via a query and persist it.
root, err := d.getFamilyRootPath(d.CloudID)
if err != nil || root == "" {
return fmt.Errorf("failed to get family root path: %w", err)
}
d.ProviderRoot = root
if len(d.Addition.RootFolderID) == 0 {
// Attempt to obtain data.path as the root via a query and persist it.
if root, err := d.getFamilyRootPath(d.CloudID); err == nil && root != "" {
d.RootFolderID = root
op.MustSaveDriverStorage(d)
}
d.RootFolderID = root
op.MustSaveDriverStorage(d)
}
_, err := d.familyGetFiles(d.RootFolderID)
_, err = d.familyGetFiles(d.RootFolderID)
if err != nil {
return err
}
Expand Down Expand Up @@ -212,7 +230,7 @@ func (d *Yun139) MakeDir(ctx context.Context, parentDir model.Obj, dirName strin
"accountType": 1,
},
"docLibName": dirName,
"path": path.Join(parentDir.GetPath(), parentDir.GetID()),
"path": d.dirPath(parentDir),
}
pathname := "/orchestration/familyCloud-rebuild/cloudCatalog/v1.0/createCloudDoc"
_, err = d.post(pathname, data, nil)
Expand All @@ -225,7 +243,7 @@ func (d *Yun139) MakeDir(ctx context.Context, parentDir model.Obj, dirName strin
},
"groupID": d.CloudID,
"parentFileId": parentDir.GetID(),
"path": path.Join(parentDir.GetPath(), parentDir.GetID()),
"path": d.dirPath(parentDir),
}
pathname := "/orchestration/group-rebuild/catalog/v1.0/createGroupCatalog"
_, err = d.post(pathname, data, nil)
Expand Down Expand Up @@ -310,9 +328,9 @@ func (d *Yun139) Move(ctx context.Context, srcObj, dstDir model.Obj) (model.Obj,
var contentList []string
var catalogList []string
if srcObj.IsDir() {
catalogList = append(catalogList, path.Join(srcObj.GetPath(), srcObj.GetID()))
catalogList = append(catalogList, d.dirPath(srcObj))
} else {
contentList = append(contentList, path.Join(srcObj.GetPath(), srcObj.GetID()))
contentList = append(contentList, d.dirPath(srcObj))
}

body := base.Json{
Expand All @@ -324,7 +342,7 @@ func (d *Yun139) Move(ctx context.Context, srcObj, dstDir model.Obj) (model.Obj,
"contentList": contentList,
"destCatalogID": dstDir.GetID(),
"destGroupID": d.CloudID,
"destPath": path.Join(dstDir.GetPath(), dstDir.GetID()),
"destPath": d.dirPath(dstDir),
"destType": 0,
"srcGroupID": d.CloudID,
"srcType": 0,
Expand Down Expand Up @@ -425,7 +443,7 @@ func (d *Yun139) Rename(ctx context.Context, srcObj model.Obj, newName string) e
},
"docLibName": newName,
"docLibraryID": srcObj.GetID(),
"path": path.Join(srcObj.GetPath(), srcObj.GetID()),
"path": d.dirPath(srcObj),
}
var resp ModifyCloudDocV2Resp
_, err = d.andAlbumRequest(pathname, data, &resp)
Expand Down Expand Up @@ -615,8 +633,22 @@ func (d *Yun139) getPartSize(size int64) int64 {
}

func (d *Yun139) Put(ctx context.Context, dstDir model.Obj, stream model.FileStreamer, up driver.UpdateProgress) error {
switch d.Addition.Type {
case MetaPersonalNew:
// PersonalNew 以及 Group/Family 在非旧流模式时走新上传路径
if d.Addition.Type == MetaPersonalNew ||
((d.isGroup() || d.isFamily()) && !d.UseOldStreamUpload) {
var createPath, getUploadUrlPath, completePath string
if d.isGroup() || d.isFamily() {
// 家庭云和共享群共用同一套新上传 API
createPath = "/dynamic/file/create"
getUploadUrlPath = "/dynamic/file/getUploadUrl"
completePath = "/dynamic/file/complete"
} else {
// MetaPersonalNew
createPath = "/file/create"
getUploadUrlPath = "/file/getUploadUrl"
completePath = "/file/complete"
}

var err error
fullHash := stream.GetHash().GetHash(utils.SHA256)
if len(fullHash) != utils.SHA256.Width {
Expand Down Expand Up @@ -671,9 +703,22 @@ func (d *Yun139) Put(ctx context.Context, dstDir model.Obj, stream model.FileStr
"type": "file",
"fileRenameMode": "auto_rename",
}
pathname := "/file/create"
// 家庭云和共享群需要额外的参数
if d.isGroup() || d.isFamily() {
if d.CloudID == "" {
return fmt.Errorf("cloud_id is required for group/family upload")
}
data["groupId"] = d.CloudID
if d.isGroup() {
data["groupType"] = 2
} else if d.isFamily() {
data["groupType"] = 1
}
data["catalogType"] = 3
data["seqNo"] = random.String(32)
}
var resp PersonalUploadResp
_, err = d.personalPost(pathname, data, &resp)
_, err = d.newPost(createPath, data, &resp)
if err != nil {
return err
}
Expand All @@ -690,10 +735,19 @@ func (d *Yun139) Put(ctx context.Context, dstDir model.Obj, stream model.FileStr
if resp.Data.PartInfos != nil {
// Progress
p := driver.NewProgress(size, up)

rateLimited := driver.NewLimitedUploadStream(ctx, stream)
ss, err := streamPkg.NewStreamSectionReader(&streamPkg.FileStream{
Ctx: ctx,
Reader: rateLimited,
Obj: &model.Object{Size: size},
}, int(partSize), &up)
if err != nil {
return err
}

// 先上传前100个分片
err = d.uploadPersonalParts(ctx, partInfos, resp.Data.PartInfos, rateLimited, p)
err = d.uploadPersonalParts(ctx, partInfos, resp.Data.PartInfos, ss, p)
if err != nil {
return err
}
Expand All @@ -711,13 +765,12 @@ func (d *Yun139) Put(ctx context.Context, dstDir model.Obj, stream model.FileStr
"accountType": 1,
},
}
pathname := "/file/getUploadUrl"
var moreresp PersonalUploadUrlResp
_, err = d.personalPost(pathname, moredata, &moreresp)
_, err = d.newPost(getUploadUrlPath, moredata, &moreresp)
if err != nil {
return err
}
err = d.uploadPersonalParts(ctx, partInfos, moreresp.Data.PartInfos, rateLimited, p)
err = d.uploadPersonalParts(ctx, partInfos, moreresp.Data.PartInfos, ss, p)
if err != nil {
return err
}
Expand All @@ -730,7 +783,11 @@ func (d *Yun139) Put(ctx context.Context, dstDir model.Obj, stream model.FileStr
"fileId": resp.Data.FileId,
"uploadId": resp.Data.UploadId,
}
_, err = d.personalPost("/file/complete", data, nil)
// 家庭云和共享群需要额外的参数
if d.isGroup() || d.isFamily() {
data["groupId"] = d.CloudID
}
_, err = d.newPost(completePath, data, nil)
if err != nil {
return err
}
Expand Down Expand Up @@ -775,11 +832,11 @@ func (d *Yun139) Put(ctx context.Context, dstDir model.Obj, stream model.FileStr
}
}
return nil
case MetaPersonal:
fallthrough
case MetaGroup:
fallthrough
case MetaFamily:
}

// 旧上传路径
switch d.Addition.Type {
case MetaPersonal, MetaGroup, MetaFamily:
// 处理冲突
// 获取文件列表
files, err := d.List(ctx, dstDir, model.ListArgs{})
Expand Down Expand Up @@ -826,11 +883,11 @@ func (d *Yun139) Put(ctx context.Context, dstDir model.Obj, stream model.FileStr
},
}
pathname := "/orchestration/personalCloud/uploadAndDownload/v1.0/pcUploadFileRequest"
if d.isFamily() || d.Addition.Type == MetaGroup {
uploadPath := path.Join(dstDir.GetPath(), dstDir.GetID())
// if dstDir is root folder
if dstDir.GetID() == d.RootFolderID {
uploadPath = d.RootPath
if d.isFamily() || d.isGroup() {
uploadPath := d.dirPath(dstDir)
Comment thread
xrgzs marked this conversation as resolved.
// 共享群的根目录上传路径为 0
if d.isGroup() && dstDir.GetID() == d.RootFolderID {
uploadPath = "0"
}
data = d.newJson(base.Json{
"fileCount": 1,
Expand Down Expand Up @@ -858,56 +915,86 @@ func (d *Yun139) Put(ctx context.Context, dstDir model.Obj, stream model.FileStr
}

size := stream.GetSize()
partSize := d.getPartSize(size)

// Progress
p := driver.NewProgress(size, up)
partSize := d.getPartSize(size)
rateLimited := driver.NewLimitedUploadStream(ctx, stream)

// StreamSectionReader for per-chunk buffering and retry
ss, err := streamPkg.NewStreamSectionReader(&streamPkg.FileStream{
Ctx: ctx,
Reader: rateLimited,
Obj: &model.Object{Size: size},
}, int(partSize), &up)
if err != nil {
return err
}

part := int64(1)
if size > partSize {
part = (size + partSize - 1) / partSize
}
rateLimited := driver.NewLimitedUploadStream(ctx, stream)
for i := int64(0); i < part; i++ {
if utils.IsCanceled(ctx) {
return ctx.Err()
}

start := i * partSize
byteSize := min(size-start, partSize)

limitReader := io.LimitReader(rateLimited, byteSize)
// Update Progress
r := io.TeeReader(limitReader, p)
req, err := http.NewRequestWithContext(ctx, http.MethodPost, resp.Data.UploadResult.RedirectionURL, r)
if err != nil {
return err
rd, getErr := ss.GetSectionReader(start, byteSize)
if getErr != nil {
return getErr
}
req.Header.Set("Content-Type", "text/plain;name="+unicode(stream.GetName()))
req.Header.Set("contentSize", strconv.FormatInt(size, 10))
req.Header.Set("range", fmt.Sprintf("bytes=%d-%d", start, start+byteSize-1))
req.Header.Set("uploadtaskID", resp.Data.UploadResult.UploadTaskID)
req.Header.Set("rangeType", "0")
req.ContentLength = byteSize

res, err := base.HttpClient.Do(req)
err = retry.Do(
func() error {
if _, err := rd.Seek(0, io.SeekStart); err != nil {
return err
}
req, reqErr := http.NewRequestWithContext(ctx, http.MethodPost, resp.Data.UploadResult.RedirectionURL,
io.TeeReader(rd, p))
if reqErr != nil {
return reqErr
}
req.Header.Set("Content-Type", "text/plain;name="+unicode(stream.GetName()))
req.Header.Set("contentSize", strconv.FormatInt(size, 10))
req.Header.Set("range", fmt.Sprintf("bytes=%d-%d", start, start+byteSize-1))
req.Header.Set("uploadtaskID", resp.Data.UploadResult.UploadTaskID)
req.Header.Set("rangeType", "0")
req.ContentLength = byteSize

res, doErr := base.HttpClient.Do(req)
if doErr != nil {
return doErr
}
defer res.Body.Close()
bodyBytes, readErr := io.ReadAll(res.Body)
if readErr != nil {
return fmt.Errorf("error reading response body: %v", readErr)
}
if res.StatusCode != http.StatusOK {
return fmt.Errorf("unexpected status code: %d, body: %s", res.StatusCode, string(bodyBytes))
}
var result InterLayerUploadResult
xmlErr := xml.Unmarshal(bodyBytes, &result)
if xmlErr != nil {
return fmt.Errorf("error parsing XML: %v", xmlErr)
}
if result.ResultCode != 0 {
return fmt.Errorf("upload failed with result code: %d, message: %s", result.ResultCode, result.Msg)
}
return nil
},
retry.Context(ctx),
retry.Attempts(3),
retry.DelayType(retry.BackOffDelay),
retry.Delay(time.Second),
)
ss.FreeSectionReader(rd)
if err != nil {
return err
}
if res.StatusCode != http.StatusOK {
res.Body.Close()
return fmt.Errorf("unexpected status code: %d", res.StatusCode)
}
bodyBytes, err := io.ReadAll(res.Body)
if err != nil {
return fmt.Errorf("error reading response body: %v", err)
}
var result InterLayerUploadResult
err = xml.Unmarshal(bodyBytes, &result)
if err != nil {
return fmt.Errorf("error parsing XML: %v", err)
}
if result.ResultCode != 0 {
return fmt.Errorf("upload failed with result code: %d, message: %s", result.ResultCode, result.Msg)
}
}
return nil
default:
Expand Down
1 change: 1 addition & 0 deletions drivers/139/meta.go
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@ type Addition struct {
CustomUploadPartSize int64 `json:"custom_upload_part_size" type:"number" default:"0" help:"0 for auto"`
ReportRealSize bool `json:"report_real_size" type:"bool" default:"true" help:"Enable to report the real file size during upload"`
UseLargeThumbnail bool `json:"use_large_thumbnail" type:"bool" default:"false" help:"Enable to use large thumbnail for images"`
UseOldStreamUpload bool `json:"use_old_stream_upload" type:"bool" default:"false" help:"Enable to use old stream upload method (not support rapid upload)"`
}

var config = driver.Config{
Expand Down
Loading