From 07e0717e1f5cd7a0b01b1dd299ebcfd196e63872 Mon Sep 17 00:00:00 2001 From: TwoOnefour Date: Tue, 27 May 2025 14:21:25 +0800 Subject: [PATCH 01/10] feat(Teldrive): Add driver Teldrive https://github.com/tgdrive/teldrive [Official api docs](https://teldrive-docs.pages.dev/docs/api) implement: * copy * move * link (302 share and local proxy) * chunk upload * rename Not yet implement: - login (scan qrcode or auth by password) - refresh token --- drivers/all.go | 1 + drivers/teldrive/driver.go | 232 +++++++++++++++ drivers/teldrive/meta.go | 27 ++ drivers/teldrive/types.go | 77 +++++ drivers/teldrive/util.go | 590 +++++++++++++++++++++++++++++++++++++ 5 files changed, 927 insertions(+) create mode 100644 drivers/teldrive/driver.go create mode 100644 drivers/teldrive/meta.go create mode 100644 drivers/teldrive/types.go create mode 100644 drivers/teldrive/util.go diff --git a/drivers/all.go b/drivers/all.go index 5b274eab0..4638d7638 100644 --- a/drivers/all.go +++ b/drivers/all.go @@ -60,6 +60,7 @@ import ( _ "github.com/OpenListTeam/OpenList/v4/drivers/smb" _ "github.com/OpenListTeam/OpenList/v4/drivers/strm" _ "github.com/OpenListTeam/OpenList/v4/drivers/teambition" + _ "github.com/OpenListTeam/OpenList/v4/drivers/teldrive" _ "github.com/OpenListTeam/OpenList/v4/drivers/terabox" _ "github.com/OpenListTeam/OpenList/v4/drivers/thunder" _ "github.com/OpenListTeam/OpenList/v4/drivers/thunder_browser" diff --git a/drivers/teldrive/driver.go b/drivers/teldrive/driver.go new file mode 100644 index 000000000..2ec3be7de --- /dev/null +++ b/drivers/teldrive/driver.go @@ -0,0 +1,232 @@ +package teldrive + +import ( + "context" + "fmt" + "github.com/OpenListTeam/OpenList/v4/drivers/base" + "github.com/OpenListTeam/OpenList/v4/internal/driver" + "github.com/OpenListTeam/OpenList/v4/internal/errs" + "github.com/OpenListTeam/OpenList/v4/internal/model" + "github.com/OpenListTeam/OpenList/v4/internal/op" + "github.com/OpenListTeam/OpenList/v4/pkg/utils" + "github.com/go-resty/resty/v2" + "github.com/google/uuid" + "math" + "net/http" + "net/url" + "strings" +) + +type Teldrive struct { + model.Storage + Addition +} + +func (d *Teldrive) Config() driver.Config { + return config +} + +func (d *Teldrive) GetAddition() driver.Additional { + return &d.Addition +} + +func (d *Teldrive) Init(ctx context.Context) error { + // TODO login / refresh token + // op.MustSaveDriverStorage(d) + if d.Cookie == "" || !strings.HasPrefix(d.Cookie, "access_token=") { + return fmt.Errorf("cookie must start with 'access_token='") + } + if d.UploadConcurrency == 0 { + d.UploadConcurrency = 4 + } + if d.ChunkSize == 0 { + d.ChunkSize = 10 + } + if d.WebdavPolicy == "native_proxy" { + d.WebProxy = true + } else { + d.WebProxy = false + } + + op.MustSaveDriverStorage(d) + return nil +} + +func (d *Teldrive) Drop(ctx context.Context) error { + return nil +} + +func (d *Teldrive) List(ctx context.Context, dir model.Obj, args model.ListArgs) ([]model.Obj, error) { + // TODO return the files list, required + // endpoint /api/files, params ->page order sort path + var listResp ListResp + params := url.Values{} + params.Set("path", dir.GetPath()) + //log.Info(dir.GetPath()) + pathname, err := utils.InjectQuery("/api/files", params) + if err != nil { + return nil, err + } + + err = d.request(http.MethodGet, pathname, nil, &listResp) + if err != nil { + return nil, err + } + + return utils.SliceConvert(listResp.Items, func(src Object) (model.Obj, error) { + return &model.Object{ + ID: src.ID, + Name: src.Name, + Size: func() int64 { + if src.Type == "folder" { + return 0 + } + return src.Size + }(), + IsFolder: src.Type == "folder", + Modified: src.UpdatedAt, + }, nil + }) +} + +func (d *Teldrive) Link(ctx context.Context, file model.Obj, args model.LinkArgs) (*model.Link, error) { + if d.WebdavPolicy != "native_proxy" { + var address string + if d.WebdavPolicy == "use_proxy_url" { + address = d.DownProxyURL + } else { + address = d.Address + } + if shareObj, err := d.getShareFileById(file.GetID()); err == nil && shareObj != nil { + return &model.Link{ + URL: address + fmt.Sprintf("/api/shares/%s/files/%s/%s", shareObj.Id, file.GetID(), file.GetName()), + }, nil + } + if err := d.createShareFile(file.GetID()); err != nil { + return nil, err + } + shareObj, err := d.getShareFileById(file.GetID()) + if err != nil { + return nil, err + } + return &model.Link{ + URL: address + fmt.Sprintf("/api/shares/%s/files/%s/%s", shareObj.Id, file.GetID(), file.GetName()), + }, nil + } + return &model.Link{ + URL: d.Address + "/api/files/" + file.GetID() + "/" + file.GetName(), + Header: http.Header{ + "Cookie": {d.Cookie}, + }, + }, nil +} + +func (d *Teldrive) MakeDir(ctx context.Context, parentDir model.Obj, dirName string) error { + return d.request(http.MethodPost, "/api/files/mkdir", func(req *resty.Request) { + req.SetBody(map[string]interface{}{ + "path": parentDir.GetPath() + "/" + dirName, + }) + }, nil) +} + +func (d *Teldrive) Move(ctx context.Context, srcObj, dstDir model.Obj) error { + body := base.Json{ + "ids": []string{srcObj.GetID()}, + "destinationParent": dstDir.GetID(), + } + return d.request(http.MethodPost, "/api/files/move", func(req *resty.Request) { + req.SetBody(body) + }, nil) +} + +func (d *Teldrive) Rename(ctx context.Context, srcObj model.Obj, newName string) error { + body := base.Json{ + "name": newName, + } + return d.request(http.MethodPatch, "/api/files/"+srcObj.GetID(), func(req *resty.Request) { + req.SetBody(body) + }, nil) +} + +func (d *Teldrive) Copy(ctx context.Context, srcObj, dstDir model.Obj) error { + copyConcurrentLimit := 4 + copyManager := NewCopyManager(ctx, copyConcurrentLimit, d) + copyManager.startWorkers() + copyManager.G.Go(func() error { + defer close(copyManager.TaskChan) + return copyManager.generateTasks(ctx, srcObj, dstDir) + }) + return copyManager.G.Wait() +} + +func (d *Teldrive) Remove(ctx context.Context, obj model.Obj) error { + body := base.Json{ + "ids": []string{obj.GetID()}, + } + return d.request(http.MethodPost, "/api/files/delete", func(req *resty.Request) { + req.SetBody(body) + }, nil) +} + +func (d *Teldrive) Put(ctx context.Context, dstDir model.Obj, file model.FileStreamer, up driver.UpdateProgress) error { + fileId := uuid.New().String() + chunkSizeInMB := d.ChunkSize + chunkSize := chunkSizeInMB * 1024 * 1024 // Convert MB to bytes + totalSize := file.GetSize() + totalParts := int(math.Ceil(float64(totalSize) / float64(chunkSize))) + retryCount := 0 + maxRetried := 3 + p := driver.NewProgress(totalSize, up) + + // delete the upload task when finished or failed + defer func() { + _ = d.request(http.MethodDelete, "/api/uploads/"+fileId, nil, nil) + }() + + if obj, err := d.getFile(dstDir.GetPath(), file.GetName(), file.IsDir()); err == nil { + if err = d.Remove(ctx, obj); err != nil { + return err + } + } + // start the upload process + if err := d.request(http.MethodGet, "/api/uploads/"+fileId, nil, nil); err != nil { + return err + } + if totalSize == 0 { + return d.touch(file.GetName(), dstDir.GetPath()) + } + + if totalParts <= 1 { + return d.doSingleUpload(ctx, dstDir, file, p, retryCount, maxRetried, totalParts, fileId) + } + + return d.doMultiUpload(ctx, dstDir, file, p, maxRetried, totalParts, chunkSize, fileId) +} + +func (d *Teldrive) GetArchiveMeta(ctx context.Context, obj model.Obj, args model.ArchiveArgs) (model.ArchiveMeta, error) { + // TODO get archive file meta-info, return errs.NotImplement to use an internal archive tool, optional + return nil, errs.NotImplement +} + +func (d *Teldrive) ListArchive(ctx context.Context, obj model.Obj, args model.ArchiveInnerArgs) ([]model.Obj, error) { + // TODO list args.InnerPath in the archive obj, return errs.NotImplement to use an internal archive tool, optional + return nil, errs.NotImplement +} + +func (d *Teldrive) Extract(ctx context.Context, obj model.Obj, args model.ArchiveInnerArgs) (*model.Link, error) { + // TODO return link of file args.InnerPath in the archive obj, return errs.NotImplement to use an internal archive tool, optional + return nil, errs.NotImplement +} + +func (d *Teldrive) ArchiveDecompress(ctx context.Context, srcObj, dstDir model.Obj, args model.ArchiveDecompressArgs) ([]model.Obj, error) { + // TODO extract args.InnerPath path in the archive srcObj to the dstDir location, optional + // a folder with the same name as the archive file needs to be created to store the extracted results if args.PutIntoNewDir + // return errs.NotImplement to use an internal archive tool + return nil, errs.NotImplement +} + +//func (d *Teldrive) Other(ctx context.Context, args model.OtherArgs) (interface{}, error) { +// return nil, errs.NotSupport +//} + +var _ driver.Driver = (*Teldrive)(nil) diff --git a/drivers/teldrive/meta.go b/drivers/teldrive/meta.go new file mode 100644 index 000000000..c5962d456 --- /dev/null +++ b/drivers/teldrive/meta.go @@ -0,0 +1,27 @@ +package teldrive + +import ( + "github.com/OpenListTeam/OpenList/v4/internal/driver" + "github.com/OpenListTeam/OpenList/v4/internal/op" +) + +type Addition struct { + // Usually one of two + driver.RootPath + // define other + Address string `json:"url" required:"true"` + ChunkSize int64 `json:"chunk_size" type:"number" default:"4" help:"Chunk size in MiB"` + Cookie string `json:"cookie" type:"string" required:"true" help:"access_token=xxx"` + UploadConcurrency int64 `json:"upload_concurrency" type:"number" default:"4" help:"Concurrency upload requests"` +} + +var config = driver.Config{ + Name: "Teldrive", + DefaultRoot: "/", +} + +func init() { + op.RegisterDriver(func() driver.Driver { + return &Teldrive{} + }) +} diff --git a/drivers/teldrive/types.go b/drivers/teldrive/types.go new file mode 100644 index 000000000..c8df4c2ff --- /dev/null +++ b/drivers/teldrive/types.go @@ -0,0 +1,77 @@ +package teldrive + +import ( + "context" + "github.com/OpenListTeam/OpenList/v4/internal/model" + "golang.org/x/sync/errgroup" + "golang.org/x/sync/semaphore" + "time" +) + +type ErrResp struct { + Code int `json:"code"` + Message string `json:"message"` +} + +type Object struct { + ID string `json:"id"` + Name string `json:"name"` + Type string `json:"type"` + MimeType string `json:"mimeType"` + Category string `json:"category,omitempty"` + ParentId string `json:"parentId"` + Size int64 `json:"size"` + Encrypted bool `json:"encrypted"` + UpdatedAt time.Time `json:"updatedAt"` +} + +type ListResp struct { + Items []Object `json:"items"` + Meta struct { + Count int `json:"count"` + TotalPages int `json:"totalPages"` + CurrentPage int `json:"currentPage"` + } `json:"meta"` +} + +type FilePart struct { + Name string `json:"name"` + PartId int `json:"partId"` + PartNo int `json:"partNo"` + ChannelId int `json:"channelId"` + Size int `json:"size"` + Encrypted bool `json:"encrypted"` + Salt string `json:"salt"` +} + +type chunkTask struct { + data []byte + chunkIdx int + fileName string +} + +type CopyManager struct { + TaskChan chan CopyTask + Sem *semaphore.Weighted + G *errgroup.Group + Ctx context.Context + d *Teldrive +} + +type CopyTask struct { + SrcObj model.Obj + DstDir model.Obj +} + +type CustomProxy struct { + model.Proxy +} + +type ShareObj struct { + Id string `json:"id"` + Protected bool `json:"protected"` + UserId int `json:"userId"` + Type string `json:"type"` + Name string `json:"name"` + ExpiresAt time.Time `json:"expiresAt"` +} diff --git a/drivers/teldrive/util.go b/drivers/teldrive/util.go new file mode 100644 index 000000000..319365fb8 --- /dev/null +++ b/drivers/teldrive/util.go @@ -0,0 +1,590 @@ +package teldrive + +import ( + "bytes" + "fmt" + "github.com/OpenListTeam/OpenList/v4/drivers/base" + "github.com/OpenListTeam/OpenList/v4/internal/driver" + "github.com/OpenListTeam/OpenList/v4/internal/model" + "github.com/OpenListTeam/OpenList/v4/pkg/utils" + "github.com/go-resty/resty/v2" + "github.com/pkg/errors" + "golang.org/x/net/context" + "golang.org/x/sync/errgroup" + "golang.org/x/sync/semaphore" + "io" + "net/http" + "sort" + "strconv" + "time" +) + +// do others that not defined in Driver interface + +func (d *Teldrive) request(method string, pathname string, callback base.ReqCallback, resp interface{}) error { + url := d.Address + pathname + req := base.RestyClient.R() + req.SetHeader("Cookie", d.Cookie) + if callback != nil { + callback(req) + } + if resp != nil { + req.SetResult(resp) + } + var e ErrResp + req.SetError(&e) + _req, err := req.Execute(method, url) + if err != nil { + return err + } + + if _req.IsError() { + return &e + } + return nil +} + +func (d *Teldrive) getFile(path, name string, isFolder bool) (model.Obj, error) { + resp := &ListResp{} + err := d.request(http.MethodGet, "/api/files", func(req *resty.Request) { + req.SetQueryParams(map[string]string{ + "path": path, + "name": name, + "type": func() string { + if isFolder { + return "folder" + } + return "file" + }(), + "operation": "find", + }) + }, resp) + if err != nil { + return nil, err + } + if len(resp.Items) == 0 { + return nil, fmt.Errorf("file not found: %s/%s", path, name) + } + obj := resp.Items[0] + return &model.Object{ + ID: obj.ID, + Name: obj.Name, + Size: obj.Size, + IsFolder: obj.Type == "folder", + }, err +} + +func (err *ErrResp) Error() string { + if err == nil { + return "" + } + + return fmt.Sprintf("[Teldrive] message:%s Error code:%d", err.Message, err.Code) +} + +// create empty file +func (d *Teldrive) touch(name, path string) error { + uploadBody := base.Json{ + "name": name, + "type": "file", + "path": path, + } + if err := d.request(http.MethodPost, "/api/files", func(req *resty.Request) { + req.SetBody(uploadBody) + }, nil); err != nil { + return err + } + + return nil +} + +func (d *Teldrive) createFileOnUploadSuccess(name, id, path string, uploadedFileParts []FilePart, totalSize int64) error { + remoteFileParts, err := d.getFilePart(id) + if err != nil { + return err + } + // check if the uploaded file parts match the remote file parts + if len(remoteFileParts) != len(uploadedFileParts) { + return fmt.Errorf("[Teldrive] file parts count mismatch: expected %d, got %d", len(uploadedFileParts), len(remoteFileParts)) + } + formatParts := make([]base.Json, 0) + for _, p := range remoteFileParts { + formatParts = append(formatParts, base.Json{ + "id": p.PartId, + "salt": p.Salt, + }) + } + uploadBody := base.Json{ + "name": name, + "type": "file", + "path": path, + "parts": formatParts, + "size": totalSize, + } + // create file here + if err := d.request(http.MethodPost, "/api/files", func(req *resty.Request) { + req.SetBody(uploadBody) + }, nil); err != nil { + return err + } + + return nil +} + +func (d *Teldrive) checkFilePartExist(fileId string, partId int) (FilePart, error) { + var uploadedParts []FilePart + var filePart FilePart + + if err := d.request(http.MethodGet, "/api/uploads/"+fileId, nil, &uploadedParts); err != nil { + return filePart, err + } + + for _, part := range uploadedParts { + if part.PartId == partId { + return part, nil + } + } + + return filePart, nil +} + +func (d *Teldrive) getFilePart(fileId string) ([]FilePart, error) { + var uploadedParts []FilePart + if err := d.request(http.MethodGet, "/api/uploads/"+fileId, nil, &uploadedParts); err != nil { + return nil, err + } + + return uploadedParts, nil +} + +func (d *Teldrive) singleUploadRequest(fileId string, callback base.ReqCallback, resp interface{}) error { + url := d.Address + "/api/uploads/" + fileId + client := resty.New().SetTimeout(0) + + ctx := context.Background() + + req := client.R(). + SetContext(ctx) + req.SetHeader("Cookie", d.Cookie) + req.SetHeader("Content-Type", "application/octet-stream") + req.SetContentLength(true) + req.AddRetryCondition(func(r *resty.Response, err error) bool { + return false + }) + if callback != nil { + callback(req) + } + if resp != nil { + req.SetResult(resp) + } + var e ErrResp + req.SetError(&e) + _req, err := req.Execute(http.MethodPost, url) + if err != nil { + return err + } + + if _req.IsError() { + return &e + } + return nil +} + +func (d *Teldrive) doSingleUpload(ctx context.Context, dstDir model.Obj, file model.FileStreamer, p *driver.Progress, + retryCount, maxRetried, totalParts int, fileId string) error { + + chunkIdx := 1 + totalSize := file.GetSize() + var fileParts []FilePart + for p.Done < p.Total { + if utils.IsCanceled(ctx) { + return ctx.Err() + } + // only one chunk, so we can use the whole file + byteData := make([]byte, totalSize) + _, err := io.ReadFull(file, byteData) + if err != io.EOF && err != nil { + return err + } + filePart := &FilePart{} + // be sure the file is uploaded, and break loop if success + for { + if err := d.singleUploadRequest(fileId, func(req *resty.Request) { + uploadParams := map[string]string{ + "partName": func() string { + digits := len(fmt.Sprintf("%d", totalParts)) + return file.GetName() + fmt.Sprintf("%0*d", digits, chunkIdx) + }(), + "partNo": strconv.Itoa(chunkIdx), + "fileName": file.GetName(), + } + req.SetQueryParams(uploadParams) + req.SetBody(driver.NewLimitedUploadStream(ctx, bytes.NewReader(byteData))) + req.SetHeader("Content-Length", strconv.Itoa(len(byteData))) + }, filePart); err != nil { + if retryCount >= maxRetried { + utils.Log.Errorf("[Teldrive] upload failed after %d retries: %s", maxRetried, err.Error()) + return err + } + if errors.Is(err, context.DeadlineExceeded) { + continue + } + retryCount++ + errorStr := fmt.Sprintf("[Teldrive] upload error: %v, retrying %d times", err, retryCount) + utils.Log.Errorf(errorStr) + time.Sleep(time.Duration(retryCount<<1) * time.Second) // Exponential backoff: 2, 4, 8, 16, ... + continue + } + break + } + if filePart.Name != "" { + fileParts = append(fileParts, *filePart) + retryCount = 0 + _, _ = p.Write(byteData) + chunkIdx++ + } + + } + + return d.createFileOnUploadSuccess(file.GetName(), fileId, dstDir.GetPath(), fileParts, totalSize) +} + +func (d *Teldrive) doMultiUpload(ctx context.Context, dstDir model.Obj, file model.FileStreamer, p *driver.Progress, + maxRetried, totalParts int, chunkSize int64, fileId string) error { + concurrent := d.UploadConcurrency + g, ctx := errgroup.WithContext(ctx) + sem := semaphore.NewWeighted(int64(concurrent)) + chunkChan := make(chan chunkTask, concurrent*2) + resultChan := make(chan FilePart, concurrent) + totalSize := file.GetSize() + g.Go(func() error { + defer close(chunkChan) + + chunkIdx := 1 + for { + select { + case <-ctx.Done(): + return ctx.Err() + default: + } + + if p.Done >= p.Total { + break + } + + byteData := make([]byte, chunkSize) + n, err := io.ReadFull(file, byteData) + if err != nil { + if errors.Is(err, io.EOF) || errors.Is(err, io.ErrUnexpectedEOF) { + if n > 0 { + // handle the last chuck + byteData = byteData[:n] + task := chunkTask{ + data: byteData, + chunkIdx: chunkIdx, + fileName: file.GetName(), + } + select { + case chunkChan <- task: + case <-ctx.Done(): + return ctx.Err() + } + } + break + } + return fmt.Errorf("read file error: %w", err) + } + + if _, err := p.Write(byteData); err != nil { + return fmt.Errorf("progress update error: %w", err) + } + + task := chunkTask{ + data: byteData, + chunkIdx: chunkIdx, + fileName: file.GetName(), + } + + select { + case chunkChan <- task: + chunkIdx++ + case <-ctx.Done(): + return ctx.Err() + } + } + return nil + }) + for i := 0; i < int(concurrent); i++ { + g.Go(func() error { + for task := range chunkChan { + if err := sem.Acquire(ctx, 1); err != nil { + return err + } + + filePart, err := d.uploadSingleChunk(ctx, fileId, task, totalParts, maxRetried) + sem.Release(1) + + if err != nil { + return fmt.Errorf("upload chunk %d failed: %w", task.chunkIdx, err) + } + + select { + case resultChan <- *filePart: + case <-ctx.Done(): + return ctx.Err() + } + } + return nil + }) + } + var fileParts []FilePart + var collectErr error + collectDone := make(chan struct{}) + + go func() { + defer close(collectDone) + fileParts = make([]FilePart, 0, totalParts) + + done := make(chan error, 1) + go func() { + done <- g.Wait() + close(resultChan) + }() + + for { + select { + case filePart, ok := <-resultChan: + if !ok { + collectErr = <-done + return + } + fileParts = append(fileParts, filePart) + case err := <-done: + collectErr = err + return + } + } + }() + + <-collectDone + + if collectErr != nil { + return fmt.Errorf("multi-upload failed: %w", collectErr) + } + sort.Slice(fileParts, func(i, j int) bool { + return fileParts[i].PartNo < fileParts[j].PartNo + }) + + return d.createFileOnUploadSuccess(file.GetName(), fileId, dstDir.GetPath(), fileParts, totalSize) +} + +func (d *Teldrive) uploadSingleChunk(ctx context.Context, fileId string, task chunkTask, totalParts, maxRetried int) (*FilePart, error) { + filePart := &FilePart{} + retryCount := 0 + + for { + select { + case <-ctx.Done(): + return nil, ctx.Err() + default: + } + + if existingPart, err := d.checkFilePartExist(fileId, task.chunkIdx); err == nil && existingPart.Name != "" { + return &existingPart, nil + } + + err := d.singleUploadRequest(fileId, func(req *resty.Request) { + uploadParams := map[string]string{ + "partName": func() string { + digits := len(fmt.Sprintf("%d", totalParts)) + return task.fileName + fmt.Sprintf("%0*d", digits, task.chunkIdx) + }(), + "partNo": strconv.Itoa(task.chunkIdx), + "fileName": task.fileName, + } + req.SetQueryParams(uploadParams) + req.SetBody(driver.NewLimitedUploadStream(ctx, bytes.NewReader(task.data))) + req.SetHeader("Content-Length", strconv.Itoa(len(task.data))) + }, filePart) + + if err == nil { + return filePart, nil + } + + if retryCount >= maxRetried { + return nil, fmt.Errorf("upload failed after %d retries: %w", maxRetried, err) + } + + if errors.Is(err, context.DeadlineExceeded) || errors.Is(err, context.Canceled) { + continue + } + + retryCount++ + utils.Log.Errorf("[Teldrive] upload error: %v, retrying %d times", err, retryCount) + + backoffDuration := time.Duration(retryCount*retryCount) * time.Second + if backoffDuration > 30*time.Second { + backoffDuration = 30 * time.Second + } + + select { + case <-time.After(backoffDuration): + case <-ctx.Done(): + return nil, ctx.Err() + } + } +} + +func (d *Teldrive) createShareFile(fileId string) error { + var errResp ErrResp + if err := d.request(http.MethodPost, "/api/files/"+fileId+"/share", func(req *resty.Request) { + req.SetBody(base.Json{ + "expiresAt": getDateTime(), + }) + }, &errResp); err != nil { + return err + } + + if errResp.Message != "" { + return &errResp + } + + return nil +} + +func (d *Teldrive) getShareFileById(fileId string) (*ShareObj, error) { + var shareObj ShareObj + if err := d.request(http.MethodGet, "/api/files/"+fileId+"/share", nil, &shareObj); err != nil { + return nil, err + } + + return &shareObj, nil +} + +func getDateTime() string { + now := time.Now().UTC() + formattedWithMs := now.Add(time.Hour * 1).Format("2006-01-02T15:04:05.000Z") + return formattedWithMs +} + +func NewCopyManager(ctx context.Context, concurrent int, d *Teldrive) *CopyManager { + g, ctx := errgroup.WithContext(ctx) + + return &CopyManager{ + TaskChan: make(chan CopyTask, concurrent*2), + Sem: semaphore.NewWeighted(int64(concurrent)), + G: g, + Ctx: ctx, + d: d, + } +} + +func (cm *CopyManager) startWorkers() { + workerCount := cap(cm.TaskChan) / 2 + for i := 0; i < workerCount; i++ { + cm.G.Go(func() error { + return cm.worker() + }) + } +} + +func (cm *CopyManager) worker() error { + for { + select { + case task, ok := <-cm.TaskChan: + if !ok { + return nil + } + + if err := cm.Sem.Acquire(cm.Ctx, 1); err != nil { + return err + } + + var err error + + err = cm.processFile(task) + + cm.Sem.Release(1) + + if err != nil { + return fmt.Errorf("task processing failed: %w", err) + } + + case <-cm.Ctx.Done(): + return cm.Ctx.Err() + } + } +} + +func (cm *CopyManager) generateTasks(ctx context.Context, srcObj, dstDir model.Obj) error { + if srcObj.IsDir() { + return cm.generateFolderTasks(ctx, srcObj, dstDir) + } else { + // add single file task directly + select { + case cm.TaskChan <- CopyTask{SrcObj: srcObj, DstDir: dstDir}: + return nil + case <-ctx.Done(): + return ctx.Err() + } + } +} + +func (cm *CopyManager) generateFolderTasks(ctx context.Context, srcDir, dstDir model.Obj) error { + objs, err := cm.d.List(ctx, srcDir, model.ListArgs{}) + if err != nil { + return fmt.Errorf("failed to list directory %s: %w", srcDir.GetPath(), err) + } + + err = cm.d.MakeDir(cm.Ctx, dstDir, srcDir.GetName()) + if err != nil || len(objs) == 0 { + return err + } + newDstDir := &model.Object{ + ID: dstDir.GetID(), + Path: dstDir.GetPath() + "/" + srcDir.GetName(), + Name: srcDir.GetName(), + IsFolder: true, + } + + for _, file := range objs { + if utils.IsCanceled(ctx) { + return ctx.Err() + } + + srcFile := &model.Object{ + ID: file.GetID(), + Path: srcDir.GetPath() + "/" + file.GetName(), + Name: file.GetName(), + IsFolder: file.IsDir(), + } + + // 递归生成任务 + if err := cm.generateTasks(ctx, srcFile, newDstDir); err != nil { + return err + } + } + + return nil +} + +func (cm *CopyManager) processFile(task CopyTask) error { + return cm.copySingleFile(cm.Ctx, task.SrcObj, task.DstDir) +} + +func (cm *CopyManager) copySingleFile(ctx context.Context, srcObj, dstDir model.Obj) error { + // `override copy mode` should delete the existing file + if obj, err := cm.d.getFile(dstDir.GetPath(), srcObj.GetName(), srcObj.IsDir()); err == nil { + if err := cm.d.Remove(ctx, obj); err != nil { + return fmt.Errorf("failed to remove existing file: %w", err) + } + } + + // Do copy + return cm.d.request(http.MethodPost, "/api/files/"+srcObj.GetID()+"/copy", func(req *resty.Request) { + req.SetBody(base.Json{ + "newName": srcObj.GetName(), + "destination": dstDir.GetPath(), + }) + }, nil) +} From f0fe503b9024f8de1db65da05c2c24d5bb0d6c43 Mon Sep 17 00:00:00 2001 From: twoonefour Date: Fri, 22 Aug 2025 15:24:55 +0800 Subject: [PATCH 02/10] fix(teldrive): force webproxy and memory optimized --- drivers/teldrive/driver.go | 57 +++++++--------- drivers/teldrive/meta.go | 1 + drivers/teldrive/types.go | 13 ++-- drivers/teldrive/util.go | 133 ++++++++++++++++++------------------- 4 files changed, 96 insertions(+), 108 deletions(-) diff --git a/drivers/teldrive/driver.go b/drivers/teldrive/driver.go index 2ec3be7de..7546bdf46 100644 --- a/drivers/teldrive/driver.go +++ b/drivers/teldrive/driver.go @@ -42,11 +42,6 @@ func (d *Teldrive) Init(ctx context.Context) error { if d.ChunkSize == 0 { d.ChunkSize = 10 } - if d.WebdavPolicy == "native_proxy" { - d.WebProxy = true - } else { - d.WebProxy = false - } op.MustSaveDriverStorage(d) return nil @@ -90,29 +85,29 @@ func (d *Teldrive) List(ctx context.Context, dir model.Obj, args model.ListArgs) } func (d *Teldrive) Link(ctx context.Context, file model.Obj, args model.LinkArgs) (*model.Link, error) { - if d.WebdavPolicy != "native_proxy" { - var address string - if d.WebdavPolicy == "use_proxy_url" { - address = d.DownProxyURL - } else { - address = d.Address - } - if shareObj, err := d.getShareFileById(file.GetID()); err == nil && shareObj != nil { - return &model.Link{ - URL: address + fmt.Sprintf("/api/shares/%s/files/%s/%s", shareObj.Id, file.GetID(), file.GetName()), - }, nil - } - if err := d.createShareFile(file.GetID()); err != nil { - return nil, err - } - shareObj, err := d.getShareFileById(file.GetID()) - if err != nil { - return nil, err - } - return &model.Link{ - URL: address + fmt.Sprintf("/api/shares/%s/files/%s/%s", shareObj.Id, file.GetID(), file.GetName()), - }, nil - } + //if d.WebdavPolicy != "native_proxy" { + // var address string + // if d.WebdavPolicy == "use_proxy_url" { + // address = d.DownProxyURL + // } else { + // address = d.Address + // } + // if shareObj, err := d.getShareFileById(file.GetID()); err == nil && shareObj != nil { + // return &model.Link{ + // URL: address + fmt.Sprintf("/api/shares/%s/files/%s/%s", shareObj.Id, file.GetID(), file.GetName()), + // }, nil + // } + // if err := d.createShareFile(file.GetID()); err != nil { + // return nil, err + // } + // shareObj, err := d.getShareFileById(file.GetID()) + // if err != nil { + // return nil, err + // } + // return &model.Link{ + // URL: address + fmt.Sprintf("/api/shares/%s/files/%s/%s", shareObj.Id, file.GetID(), file.GetName()), + // }, nil + //} return &model.Link{ URL: d.Address + "/api/files/" + file.GetID() + "/" + file.GetName(), Header: http.Header{ @@ -174,9 +169,7 @@ func (d *Teldrive) Put(ctx context.Context, dstDir model.Obj, file model.FileStr chunkSize := chunkSizeInMB * 1024 * 1024 // Convert MB to bytes totalSize := file.GetSize() totalParts := int(math.Ceil(float64(totalSize) / float64(chunkSize))) - retryCount := 0 maxRetried := 3 - p := driver.NewProgress(totalSize, up) // delete the upload task when finished or failed defer func() { @@ -197,10 +190,10 @@ func (d *Teldrive) Put(ctx context.Context, dstDir model.Obj, file model.FileStr } if totalParts <= 1 { - return d.doSingleUpload(ctx, dstDir, file, p, retryCount, maxRetried, totalParts, fileId) + return d.doSingleUpload(ctx, dstDir, file, up, totalParts, chunkSize, fileId) } - return d.doMultiUpload(ctx, dstDir, file, p, maxRetried, totalParts, chunkSize, fileId) + return d.doMultiUpload(ctx, dstDir, file, up, maxRetried, totalParts, chunkSize, fileId) } func (d *Teldrive) GetArchiveMeta(ctx context.Context, obj model.Obj, args model.ArchiveArgs) (model.ArchiveMeta, error) { diff --git a/drivers/teldrive/meta.go b/drivers/teldrive/meta.go index c5962d456..cd25a7d18 100644 --- a/drivers/teldrive/meta.go +++ b/drivers/teldrive/meta.go @@ -18,6 +18,7 @@ type Addition struct { var config = driver.Config{ Name: "Teldrive", DefaultRoot: "/", + OnlyProxy: true, } func init() { diff --git a/drivers/teldrive/types.go b/drivers/teldrive/types.go index c8df4c2ff..c7ef7603c 100644 --- a/drivers/teldrive/types.go +++ b/drivers/teldrive/types.go @@ -3,6 +3,7 @@ package teldrive import ( "context" "github.com/OpenListTeam/OpenList/v4/internal/model" + "github.com/OpenListTeam/OpenList/v4/internal/stream" "golang.org/x/sync/errgroup" "golang.org/x/sync/semaphore" "time" @@ -45,9 +46,11 @@ type FilePart struct { } type chunkTask struct { - data []byte - chunkIdx int - fileName string + chunkIdx int + fileName string + chunkSize int64 + reader *stream.SectionReader + ss *stream.StreamSectionReader } type CopyManager struct { @@ -63,10 +66,6 @@ type CopyTask struct { DstDir model.Obj } -type CustomProxy struct { - model.Proxy -} - type ShareObj struct { Id string `json:"id"` Protected bool `json:"protected"` diff --git a/drivers/teldrive/util.go b/drivers/teldrive/util.go index 319365fb8..4c388b767 100644 --- a/drivers/teldrive/util.go +++ b/drivers/teldrive/util.go @@ -1,12 +1,13 @@ package teldrive import ( - "bytes" "fmt" "github.com/OpenListTeam/OpenList/v4/drivers/base" "github.com/OpenListTeam/OpenList/v4/internal/driver" "github.com/OpenListTeam/OpenList/v4/internal/model" + "github.com/OpenListTeam/OpenList/v4/internal/stream" "github.com/OpenListTeam/OpenList/v4/pkg/utils" + "github.com/avast/retry-go" "github.com/go-resty/resty/v2" "github.com/pkg/errors" "golang.org/x/net/context" @@ -16,6 +17,7 @@ import ( "net/http" "sort" "strconv" + "sync" "time" ) @@ -190,58 +192,62 @@ func (d *Teldrive) singleUploadRequest(fileId string, callback base.ReqCallback, return nil } -func (d *Teldrive) doSingleUpload(ctx context.Context, dstDir model.Obj, file model.FileStreamer, p *driver.Progress, - retryCount, maxRetried, totalParts int, fileId string) error { +func (d *Teldrive) doSingleUpload(ctx context.Context, dstDir model.Obj, file model.FileStreamer, up model.UpdateProgress, + totalParts int, chunkSize int64, fileId string) error { - chunkIdx := 1 totalSize := file.GetSize() var fileParts []FilePart - for p.Done < p.Total { + var uploaded int64 = 0 + ss, err := stream.NewStreamSectionReader(file, int(totalSize), &up) + if err != nil { + return err + } + + for uploaded < totalSize { if utils.IsCanceled(ctx) { return ctx.Err() } - // only one chunk, so we can use the whole file - byteData := make([]byte, totalSize) - _, err := io.ReadFull(file, byteData) - if err != io.EOF && err != nil { + curChunkSize := min(totalSize-uploaded, chunkSize) + rd, err := ss.GetSectionReader(uploaded, curChunkSize) + if err != nil { return err } filePart := &FilePart{} - // be sure the file is uploaded, and break loop if success - for { + if err := retry.Do(func() error { + + if _, err := rd.Seek(0, io.SeekStart); err != nil { + return err + } + if err := d.singleUploadRequest(fileId, func(req *resty.Request) { uploadParams := map[string]string{ "partName": func() string { digits := len(fmt.Sprintf("%d", totalParts)) - return file.GetName() + fmt.Sprintf("%0*d", digits, chunkIdx) + return file.GetName() + fmt.Sprintf(".%0*d", digits, 1) }(), - "partNo": strconv.Itoa(chunkIdx), + "partNo": strconv.Itoa(1), "fileName": file.GetName(), } req.SetQueryParams(uploadParams) - req.SetBody(driver.NewLimitedUploadStream(ctx, bytes.NewReader(byteData))) - req.SetHeader("Content-Length", strconv.Itoa(len(byteData))) + req.SetBody(driver.NewLimitedUploadStream(ctx, rd)) + req.SetHeader("Content-Length", strconv.FormatInt(curChunkSize, 10)) }, filePart); err != nil { - if retryCount >= maxRetried { - utils.Log.Errorf("[Teldrive] upload failed after %d retries: %s", maxRetried, err.Error()) - return err - } - if errors.Is(err, context.DeadlineExceeded) { - continue - } - retryCount++ - errorStr := fmt.Sprintf("[Teldrive] upload error: %v, retrying %d times", err, retryCount) - utils.Log.Errorf(errorStr) - time.Sleep(time.Duration(retryCount<<1) * time.Second) // Exponential backoff: 2, 4, 8, 16, ... - continue + return err } - break + + return nil + }, + retry.Attempts(3), + retry.DelayType(retry.BackOffDelay), + retry.Delay(time.Second)); err != nil { + return err } + if filePart.Name != "" { fileParts = append(fileParts, *filePart) - retryCount = 0 - _, _ = p.Write(byteData) - chunkIdx++ + uploaded += curChunkSize + up(float64(uploaded) / float64(totalSize)) + ss.FreeSectionReader(rd) } } @@ -249,62 +255,50 @@ func (d *Teldrive) doSingleUpload(ctx context.Context, dstDir model.Obj, file mo return d.createFileOnUploadSuccess(file.GetName(), fileId, dstDir.GetPath(), fileParts, totalSize) } -func (d *Teldrive) doMultiUpload(ctx context.Context, dstDir model.Obj, file model.FileStreamer, p *driver.Progress, +func (d *Teldrive) doMultiUpload(ctx context.Context, dstDir model.Obj, file model.FileStreamer, up model.UpdateProgress, maxRetried, totalParts int, chunkSize int64, fileId string) error { + concurrent := d.UploadConcurrency g, ctx := errgroup.WithContext(ctx) sem := semaphore.NewWeighted(int64(concurrent)) chunkChan := make(chan chunkTask, concurrent*2) resultChan := make(chan FilePart, concurrent) totalSize := file.GetSize() + + ss, err := stream.NewStreamSectionReader(file, int(totalSize), &up) + if err != nil { + return err + } + ssLock := sync.Mutex{} g.Go(func() error { defer close(chunkChan) - chunkIdx := 1 - for { + chunkIdx := 0 + for chunkIdx < totalParts { select { case <-ctx.Done(): return ctx.Err() default: } - if p.Done >= p.Total { - break - } + offset := int64(chunkIdx) * chunkSize + curChunkSize := min(totalSize-offset, chunkSize) - byteData := make([]byte, chunkSize) - n, err := io.ReadFull(file, byteData) - if err != nil { - if errors.Is(err, io.EOF) || errors.Is(err, io.ErrUnexpectedEOF) { - if n > 0 { - // handle the last chuck - byteData = byteData[:n] - task := chunkTask{ - data: byteData, - chunkIdx: chunkIdx, - fileName: file.GetName(), - } - select { - case chunkChan <- task: - case <-ctx.Done(): - return ctx.Err() - } - } - break - } - return fmt.Errorf("read file error: %w", err) - } + ssLock.Lock() + reader, err := ss.GetSectionReader(offset, curChunkSize) + ssLock.Unlock() - if _, err := p.Write(byteData); err != nil { - return fmt.Errorf("progress update error: %w", err) + if err != nil { + return err } - task := chunkTask{ - data: byteData, - chunkIdx: chunkIdx, - fileName: file.GetName(), + chunkIdx: chunkIdx + 1, + chunkSize: curChunkSize, + fileName: file.GetName(), + reader: reader, + ss: ss, } - + // freeSectionReader will be called in d.uploadSingleChunk select { case chunkChan <- task: chunkIdx++ @@ -381,6 +375,7 @@ func (d *Teldrive) doMultiUpload(ctx context.Context, dstDir model.Obj, file mod func (d *Teldrive) uploadSingleChunk(ctx context.Context, fileId string, task chunkTask, totalParts, maxRetried int) (*FilePart, error) { filePart := &FilePart{} retryCount := 0 + defer task.ss.FreeSectionReader(task.reader) for { select { @@ -397,14 +392,14 @@ func (d *Teldrive) uploadSingleChunk(ctx context.Context, fileId string, task ch uploadParams := map[string]string{ "partName": func() string { digits := len(fmt.Sprintf("%d", totalParts)) - return task.fileName + fmt.Sprintf("%0*d", digits, task.chunkIdx) + return task.fileName + fmt.Sprintf(".%0*d", digits, task.chunkIdx) }(), "partNo": strconv.Itoa(task.chunkIdx), "fileName": task.fileName, } req.SetQueryParams(uploadParams) - req.SetBody(driver.NewLimitedUploadStream(ctx, bytes.NewReader(task.data))) - req.SetHeader("Content-Length", strconv.Itoa(len(task.data))) + req.SetBody(driver.NewLimitedUploadStream(ctx, task.reader)) + req.SetHeader("Content-Length", strconv.Itoa(int(task.chunkSize))) }, filePart) if err == nil { From 57cdaea1f17507a77ad0e0039c1b722c481d4d23 Mon Sep 17 00:00:00 2001 From: MadDogOwner Date: Sun, 24 Aug 2025 02:36:17 +0800 Subject: [PATCH 03/10] chore(teldrive): go fmt --- drivers/teldrive/driver.go | 9 +++++---- drivers/teldrive/types.go | 3 ++- drivers/teldrive/util.go | 13 +++++++------ 3 files changed, 14 insertions(+), 11 deletions(-) diff --git a/drivers/teldrive/driver.go b/drivers/teldrive/driver.go index 7546bdf46..e999a7907 100644 --- a/drivers/teldrive/driver.go +++ b/drivers/teldrive/driver.go @@ -3,6 +3,11 @@ package teldrive import ( "context" "fmt" + "math" + "net/http" + "net/url" + "strings" + "github.com/OpenListTeam/OpenList/v4/drivers/base" "github.com/OpenListTeam/OpenList/v4/internal/driver" "github.com/OpenListTeam/OpenList/v4/internal/errs" @@ -11,10 +16,6 @@ import ( "github.com/OpenListTeam/OpenList/v4/pkg/utils" "github.com/go-resty/resty/v2" "github.com/google/uuid" - "math" - "net/http" - "net/url" - "strings" ) type Teldrive struct { diff --git a/drivers/teldrive/types.go b/drivers/teldrive/types.go index c7ef7603c..f6399e06a 100644 --- a/drivers/teldrive/types.go +++ b/drivers/teldrive/types.go @@ -2,11 +2,12 @@ package teldrive import ( "context" + "time" + "github.com/OpenListTeam/OpenList/v4/internal/model" "github.com/OpenListTeam/OpenList/v4/internal/stream" "golang.org/x/sync/errgroup" "golang.org/x/sync/semaphore" - "time" ) type ErrResp struct { diff --git a/drivers/teldrive/util.go b/drivers/teldrive/util.go index 4c388b767..1d2ac2fe9 100644 --- a/drivers/teldrive/util.go +++ b/drivers/teldrive/util.go @@ -2,6 +2,13 @@ package teldrive import ( "fmt" + "io" + "net/http" + "sort" + "strconv" + "sync" + "time" + "github.com/OpenListTeam/OpenList/v4/drivers/base" "github.com/OpenListTeam/OpenList/v4/internal/driver" "github.com/OpenListTeam/OpenList/v4/internal/model" @@ -13,12 +20,6 @@ import ( "golang.org/x/net/context" "golang.org/x/sync/errgroup" "golang.org/x/sync/semaphore" - "io" - "net/http" - "sort" - "strconv" - "sync" - "time" ) // do others that not defined in Driver interface From d5f8ee999b2aaa72963c4be1edf1629fb130c8f8 Mon Sep 17 00:00:00 2001 From: MadDogOwner Date: Sun, 24 Aug 2025 03:07:01 +0800 Subject: [PATCH 04/10] chore(teldrive): remove TODO --- drivers/teldrive/driver.go | 3 --- 1 file changed, 3 deletions(-) diff --git a/drivers/teldrive/driver.go b/drivers/teldrive/driver.go index e999a7907..38d40bd69 100644 --- a/drivers/teldrive/driver.go +++ b/drivers/teldrive/driver.go @@ -32,8 +32,6 @@ func (d *Teldrive) GetAddition() driver.Additional { } func (d *Teldrive) Init(ctx context.Context) error { - // TODO login / refresh token - // op.MustSaveDriverStorage(d) if d.Cookie == "" || !strings.HasPrefix(d.Cookie, "access_token=") { return fmt.Errorf("cookie must start with 'access_token='") } @@ -53,7 +51,6 @@ func (d *Teldrive) Drop(ctx context.Context) error { } func (d *Teldrive) List(ctx context.Context, dir model.Obj, args model.ListArgs) ([]model.Obj, error) { - // TODO return the files list, required // endpoint /api/files, params ->page order sort path var listResp ListResp params := url.Values{} From 329abbbc4a1120791d807aaf8905d96133aacf3d Mon Sep 17 00:00:00 2001 From: MadDogOwner Date: Sun, 24 Aug 2025 03:18:19 +0800 Subject: [PATCH 05/10] chore(teldrive): organize code --- drivers/teldrive/copy.go | 136 +++++++++++ drivers/teldrive/upload.go | 369 ++++++++++++++++++++++++++++ drivers/teldrive/util.go | 480 ------------------------------------- 3 files changed, 505 insertions(+), 480 deletions(-) create mode 100644 drivers/teldrive/copy.go create mode 100644 drivers/teldrive/upload.go diff --git a/drivers/teldrive/copy.go b/drivers/teldrive/copy.go new file mode 100644 index 000000000..6495d702d --- /dev/null +++ b/drivers/teldrive/copy.go @@ -0,0 +1,136 @@ +package teldrive + +import ( + "fmt" + "net/http" + + "github.com/OpenListTeam/OpenList/v4/drivers/base" + "github.com/OpenListTeam/OpenList/v4/internal/model" + "github.com/OpenListTeam/OpenList/v4/pkg/utils" + "github.com/go-resty/resty/v2" + "golang.org/x/net/context" + "golang.org/x/sync/errgroup" + "golang.org/x/sync/semaphore" +) + +func NewCopyManager(ctx context.Context, concurrent int, d *Teldrive) *CopyManager { + g, ctx := errgroup.WithContext(ctx) + + return &CopyManager{ + TaskChan: make(chan CopyTask, concurrent*2), + Sem: semaphore.NewWeighted(int64(concurrent)), + G: g, + Ctx: ctx, + d: d, + } +} + +func (cm *CopyManager) startWorkers() { + workerCount := cap(cm.TaskChan) / 2 + for i := 0; i < workerCount; i++ { + cm.G.Go(func() error { + return cm.worker() + }) + } +} + +func (cm *CopyManager) worker() error { + for { + select { + case task, ok := <-cm.TaskChan: + if !ok { + return nil + } + + if err := cm.Sem.Acquire(cm.Ctx, 1); err != nil { + return err + } + + var err error + + err = cm.processFile(task) + + cm.Sem.Release(1) + + if err != nil { + return fmt.Errorf("task processing failed: %w", err) + } + + case <-cm.Ctx.Done(): + return cm.Ctx.Err() + } + } +} + +func (cm *CopyManager) generateTasks(ctx context.Context, srcObj, dstDir model.Obj) error { + if srcObj.IsDir() { + return cm.generateFolderTasks(ctx, srcObj, dstDir) + } else { + // add single file task directly + select { + case cm.TaskChan <- CopyTask{SrcObj: srcObj, DstDir: dstDir}: + return nil + case <-ctx.Done(): + return ctx.Err() + } + } +} + +func (cm *CopyManager) generateFolderTasks(ctx context.Context, srcDir, dstDir model.Obj) error { + objs, err := cm.d.List(ctx, srcDir, model.ListArgs{}) + if err != nil { + return fmt.Errorf("failed to list directory %s: %w", srcDir.GetPath(), err) + } + + err = cm.d.MakeDir(cm.Ctx, dstDir, srcDir.GetName()) + if err != nil || len(objs) == 0 { + return err + } + newDstDir := &model.Object{ + ID: dstDir.GetID(), + Path: dstDir.GetPath() + "/" + srcDir.GetName(), + Name: srcDir.GetName(), + IsFolder: true, + } + + for _, file := range objs { + if utils.IsCanceled(ctx) { + return ctx.Err() + } + + srcFile := &model.Object{ + ID: file.GetID(), + Path: srcDir.GetPath() + "/" + file.GetName(), + Name: file.GetName(), + IsFolder: file.IsDir(), + } + + // 递归生成任务 + if err := cm.generateTasks(ctx, srcFile, newDstDir); err != nil { + return err + } + } + + return nil +} + +func (cm *CopyManager) processFile(task CopyTask) error { + return cm.copySingleFile(cm.Ctx, task.SrcObj, task.DstDir) +} + +func (cm *CopyManager) copySingleFile(ctx context.Context, srcObj, dstDir model.Obj) error { + // `override copy mode` should delete the existing file + if obj, err := cm.d.getFile(dstDir.GetPath(), srcObj.GetName(), srcObj.IsDir()); err == nil { + if err := cm.d.Remove(ctx, obj); err != nil { + return fmt.Errorf("failed to remove existing file: %w", err) + } + } + + // Do copy + return cm.d.request(http.MethodPost, "/api/files/"+srcObj.GetID()+"/copy", func(req *resty.Request) { + req.SetBody(base.Json{ + "newName": srcObj.GetName(), + "destination": dstDir.GetPath(), + }) + }, nil) +} diff --git a/drivers/teldrive/upload.go b/drivers/teldrive/upload.go new file mode 100644 index 000000000..7efb995b6 --- /dev/null +++ b/drivers/teldrive/upload.go @@ -0,0 +1,369 @@ +package teldrive + +import ( + "fmt" + "io" + "net/http" + "sort" + "strconv" + "sync" + "time" + + "github.com/OpenListTeam/OpenList/v4/drivers/base" + "github.com/OpenListTeam/OpenList/v4/internal/driver" + "github.com/OpenListTeam/OpenList/v4/internal/model" + "github.com/OpenListTeam/OpenList/v4/internal/stream" + "github.com/OpenListTeam/OpenList/v4/pkg/utils" + "github.com/avast/retry-go" + "github.com/go-resty/resty/v2" + "github.com/pkg/errors" + "golang.org/x/net/context" + "golang.org/x/sync/errgroup" + "golang.org/x/sync/semaphore" +) + +// create empty file +func (d *Teldrive) touch(name, path string) error { + uploadBody := base.Json{ + "name": name, + "type": "file", + "path": path, + } + if err := d.request(http.MethodPost, "/api/files", func(req *resty.Request) { + req.SetBody(uploadBody) + }, nil); err != nil { + return err + } + + return nil +} + +func (d *Teldrive) createFileOnUploadSuccess(name, id, path string, uploadedFileParts []FilePart, totalSize int64) error { + remoteFileParts, err := d.getFilePart(id) + if err != nil { + return err + } + // check if the uploaded file parts match the remote file parts + if len(remoteFileParts) != len(uploadedFileParts) { + return fmt.Errorf("[Teldrive] file parts count mismatch: expected %d, got %d", len(uploadedFileParts), len(remoteFileParts)) + } + formatParts := make([]base.Json, 0) + for _, p := range remoteFileParts { + formatParts = append(formatParts, base.Json{ + "id": p.PartId, + "salt": p.Salt, + }) + } + uploadBody := base.Json{ + "name": name, + "type": "file", + "path": path, + "parts": formatParts, + "size": totalSize, + } + // create file here + if err := d.request(http.MethodPost, "/api/files", func(req *resty.Request) { + req.SetBody(uploadBody) + }, nil); err != nil { + return err + } + + return nil +} + +func (d *Teldrive) checkFilePartExist(fileId string, partId int) (FilePart, error) { + var uploadedParts []FilePart + var filePart FilePart + + if err := d.request(http.MethodGet, "/api/uploads/"+fileId, nil, &uploadedParts); err != nil { + return filePart, err + } + + for _, part := range uploadedParts { + if part.PartId == partId { + return part, nil + } + } + + return filePart, nil +} + +func (d *Teldrive) getFilePart(fileId string) ([]FilePart, error) { + var uploadedParts []FilePart + if err := d.request(http.MethodGet, "/api/uploads/"+fileId, nil, &uploadedParts); err != nil { + return nil, err + } + + return uploadedParts, nil +} + +func (d *Teldrive) singleUploadRequest(fileId string, callback base.ReqCallback, resp interface{}) error { + url := d.Address + "/api/uploads/" + fileId + client := resty.New().SetTimeout(0) + + ctx := context.Background() + + req := client.R(). + SetContext(ctx) + req.SetHeader("Cookie", d.Cookie) + req.SetHeader("Content-Type", "application/octet-stream") + req.SetContentLength(true) + req.AddRetryCondition(func(r *resty.Response, err error) bool { + return false + }) + if callback != nil { + callback(req) + } + if resp != nil { + req.SetResult(resp) + } + var e ErrResp + req.SetError(&e) + _req, err := req.Execute(http.MethodPost, url) + if err != nil { + return err + } + + if _req.IsError() { + return &e + } + return nil +} + +func (d *Teldrive) doSingleUpload(ctx context.Context, dstDir model.Obj, file model.FileStreamer, up model.UpdateProgress, + totalParts int, chunkSize int64, fileId string) error { + + totalSize := file.GetSize() + var fileParts []FilePart + var uploaded int64 = 0 + ss, err := stream.NewStreamSectionReader(file, int(totalSize), &up) + if err != nil { + return err + } + + for uploaded < totalSize { + if utils.IsCanceled(ctx) { + return ctx.Err() + } + curChunkSize := min(totalSize-uploaded, chunkSize) + rd, err := ss.GetSectionReader(uploaded, curChunkSize) + if err != nil { + return err + } + filePart := &FilePart{} + if err := retry.Do(func() error { + + if _, err := rd.Seek(0, io.SeekStart); err != nil { + return err + } + + if err := d.singleUploadRequest(fileId, func(req *resty.Request) { + uploadParams := map[string]string{ + "partName": func() string { + digits := len(fmt.Sprintf("%d", totalParts)) + return file.GetName() + fmt.Sprintf(".%0*d", digits, 1) + }(), + "partNo": strconv.Itoa(1), + "fileName": file.GetName(), + } + req.SetQueryParams(uploadParams) + req.SetBody(driver.NewLimitedUploadStream(ctx, rd)) + req.SetHeader("Content-Length", strconv.FormatInt(curChunkSize, 10)) + }, filePart); err != nil { + return err + } + + return nil + }, + retry.Attempts(3), + retry.DelayType(retry.BackOffDelay), + retry.Delay(time.Second)); err != nil { + return err + } + + if filePart.Name != "" { + fileParts = append(fileParts, *filePart) + uploaded += curChunkSize + up(float64(uploaded) / float64(totalSize)) + ss.FreeSectionReader(rd) + } + + } + + return d.createFileOnUploadSuccess(file.GetName(), fileId, dstDir.GetPath(), fileParts, totalSize) +} + +func (d *Teldrive) doMultiUpload(ctx context.Context, dstDir model.Obj, file model.FileStreamer, up model.UpdateProgress, + maxRetried, totalParts int, chunkSize int64, fileId string) error { + + concurrent := d.UploadConcurrency + g, ctx := errgroup.WithContext(ctx) + sem := semaphore.NewWeighted(int64(concurrent)) + chunkChan := make(chan chunkTask, concurrent*2) + resultChan := make(chan FilePart, concurrent) + totalSize := file.GetSize() + + ss, err := stream.NewStreamSectionReader(file, int(totalSize), &up) + if err != nil { + return err + } + ssLock := sync.Mutex{} + g.Go(func() error { + defer close(chunkChan) + + chunkIdx := 0 + for chunkIdx < totalParts { + select { + case <-ctx.Done(): + return ctx.Err() + default: + } + + offset := int64(chunkIdx) * chunkSize + curChunkSize := min(totalSize-offset, chunkSize) + + ssLock.Lock() + reader, err := ss.GetSectionReader(offset, curChunkSize) + ssLock.Unlock() + + if err != nil { + return err + } + task := chunkTask{ + chunkIdx: chunkIdx + 1, + chunkSize: curChunkSize, + fileName: file.GetName(), + reader: reader, + ss: ss, + } + // freeSectionReader will be called in d.uploadSingleChunk + select { + case chunkChan <- task: + chunkIdx++ + case <-ctx.Done(): + return ctx.Err() + } + } + return nil + }) + for i := 0; i < int(concurrent); i++ { + g.Go(func() error { + for task := range chunkChan { + if err := sem.Acquire(ctx, 1); err != nil { + return err + } + + filePart, err := d.uploadSingleChunk(ctx, fileId, task, totalParts, maxRetried) + sem.Release(1) + + if err != nil { + return fmt.Errorf("upload chunk %d failed: %w", task.chunkIdx, err) + } + + select { + case resultChan <- *filePart: + case <-ctx.Done(): + return ctx.Err() + } + } + return nil + }) + } + var fileParts []FilePart + var collectErr error + collectDone := make(chan struct{}) + + go func() { + defer close(collectDone) + fileParts = make([]FilePart, 0, totalParts) + + done := make(chan error, 1) + go func() { + done <- g.Wait() + close(resultChan) + }() + + for { + select { + case filePart, ok := <-resultChan: + if !ok { + collectErr = <-done + return + } + fileParts = append(fileParts, filePart) + case err := <-done: + collectErr = err + return + } + } + }() + + <-collectDone + + if collectErr != nil { + return fmt.Errorf("multi-upload failed: %w", collectErr) + } + sort.Slice(fileParts, func(i, j int) bool { + return fileParts[i].PartNo < fileParts[j].PartNo + }) + + return d.createFileOnUploadSuccess(file.GetName(), fileId, dstDir.GetPath(), fileParts, totalSize) +} + +func (d *Teldrive) uploadSingleChunk(ctx context.Context, fileId string, task chunkTask, totalParts, maxRetried int) (*FilePart, error) { + filePart := &FilePart{} + retryCount := 0 + defer task.ss.FreeSectionReader(task.reader) + + for { + select { + case <-ctx.Done(): + return nil, ctx.Err() + default: + } + + if existingPart, err := d.checkFilePartExist(fileId, task.chunkIdx); err == nil && existingPart.Name != "" { + return &existingPart, nil + } + + err := d.singleUploadRequest(fileId, func(req *resty.Request) { + uploadParams := map[string]string{ + "partName": func() string { + digits := len(fmt.Sprintf("%d", totalParts)) + return task.fileName + fmt.Sprintf(".%0*d", digits, task.chunkIdx) + }(), + "partNo": strconv.Itoa(task.chunkIdx), + "fileName": task.fileName, + } + req.SetQueryParams(uploadParams) + req.SetBody(driver.NewLimitedUploadStream(ctx, task.reader)) + req.SetHeader("Content-Length", strconv.Itoa(int(task.chunkSize))) + }, filePart) + + if err == nil { + return filePart, nil + } + + if retryCount >= maxRetried { + return nil, fmt.Errorf("upload failed after %d retries: %w", maxRetried, err) + } + + if errors.Is(err, context.DeadlineExceeded) || errors.Is(err, context.Canceled) { + continue + } + + retryCount++ + utils.Log.Errorf("[Teldrive] upload error: %v, retrying %d times", err, retryCount) + + backoffDuration := time.Duration(retryCount*retryCount) * time.Second + if backoffDuration > 30*time.Second { + backoffDuration = 30 * time.Second + } + + select { + case <-time.After(backoffDuration): + case <-ctx.Done(): + return nil, ctx.Err() + } + } +} diff --git a/drivers/teldrive/util.go b/drivers/teldrive/util.go index 1d2ac2fe9..bda9d4506 100644 --- a/drivers/teldrive/util.go +++ b/drivers/teldrive/util.go @@ -2,24 +2,12 @@ package teldrive import ( "fmt" - "io" "net/http" - "sort" - "strconv" - "sync" "time" "github.com/OpenListTeam/OpenList/v4/drivers/base" - "github.com/OpenListTeam/OpenList/v4/internal/driver" "github.com/OpenListTeam/OpenList/v4/internal/model" - "github.com/OpenListTeam/OpenList/v4/internal/stream" - "github.com/OpenListTeam/OpenList/v4/pkg/utils" - "github.com/avast/retry-go" "github.com/go-resty/resty/v2" - "github.com/pkg/errors" - "golang.org/x/net/context" - "golang.org/x/sync/errgroup" - "golang.org/x/sync/semaphore" ) // do others that not defined in Driver interface @@ -85,352 +73,6 @@ func (err *ErrResp) Error() string { return fmt.Sprintf("[Teldrive] message:%s Error code:%d", err.Message, err.Code) } -// create empty file -func (d *Teldrive) touch(name, path string) error { - uploadBody := base.Json{ - "name": name, - "type": "file", - "path": path, - } - if err := d.request(http.MethodPost, "/api/files", func(req *resty.Request) { - req.SetBody(uploadBody) - }, nil); err != nil { - return err - } - - return nil -} - -func (d *Teldrive) createFileOnUploadSuccess(name, id, path string, uploadedFileParts []FilePart, totalSize int64) error { - remoteFileParts, err := d.getFilePart(id) - if err != nil { - return err - } - // check if the uploaded file parts match the remote file parts - if len(remoteFileParts) != len(uploadedFileParts) { - return fmt.Errorf("[Teldrive] file parts count mismatch: expected %d, got %d", len(uploadedFileParts), len(remoteFileParts)) - } - formatParts := make([]base.Json, 0) - for _, p := range remoteFileParts { - formatParts = append(formatParts, base.Json{ - "id": p.PartId, - "salt": p.Salt, - }) - } - uploadBody := base.Json{ - "name": name, - "type": "file", - "path": path, - "parts": formatParts, - "size": totalSize, - } - // create file here - if err := d.request(http.MethodPost, "/api/files", func(req *resty.Request) { - req.SetBody(uploadBody) - }, nil); err != nil { - return err - } - - return nil -} - -func (d *Teldrive) checkFilePartExist(fileId string, partId int) (FilePart, error) { - var uploadedParts []FilePart - var filePart FilePart - - if err := d.request(http.MethodGet, "/api/uploads/"+fileId, nil, &uploadedParts); err != nil { - return filePart, err - } - - for _, part := range uploadedParts { - if part.PartId == partId { - return part, nil - } - } - - return filePart, nil -} - -func (d *Teldrive) getFilePart(fileId string) ([]FilePart, error) { - var uploadedParts []FilePart - if err := d.request(http.MethodGet, "/api/uploads/"+fileId, nil, &uploadedParts); err != nil { - return nil, err - } - - return uploadedParts, nil -} - -func (d *Teldrive) singleUploadRequest(fileId string, callback base.ReqCallback, resp interface{}) error { - url := d.Address + "/api/uploads/" + fileId - client := resty.New().SetTimeout(0) - - ctx := context.Background() - - req := client.R(). - SetContext(ctx) - req.SetHeader("Cookie", d.Cookie) - req.SetHeader("Content-Type", "application/octet-stream") - req.SetContentLength(true) - req.AddRetryCondition(func(r *resty.Response, err error) bool { - return false - }) - if callback != nil { - callback(req) - } - if resp != nil { - req.SetResult(resp) - } - var e ErrResp - req.SetError(&e) - _req, err := req.Execute(http.MethodPost, url) - if err != nil { - return err - } - - if _req.IsError() { - return &e - } - return nil -} - -func (d *Teldrive) doSingleUpload(ctx context.Context, dstDir model.Obj, file model.FileStreamer, up model.UpdateProgress, - totalParts int, chunkSize int64, fileId string) error { - - totalSize := file.GetSize() - var fileParts []FilePart - var uploaded int64 = 0 - ss, err := stream.NewStreamSectionReader(file, int(totalSize), &up) - if err != nil { - return err - } - - for uploaded < totalSize { - if utils.IsCanceled(ctx) { - return ctx.Err() - } - curChunkSize := min(totalSize-uploaded, chunkSize) - rd, err := ss.GetSectionReader(uploaded, curChunkSize) - if err != nil { - return err - } - filePart := &FilePart{} - if err := retry.Do(func() error { - - if _, err := rd.Seek(0, io.SeekStart); err != nil { - return err - } - - if err := d.singleUploadRequest(fileId, func(req *resty.Request) { - uploadParams := map[string]string{ - "partName": func() string { - digits := len(fmt.Sprintf("%d", totalParts)) - return file.GetName() + fmt.Sprintf(".%0*d", digits, 1) - }(), - "partNo": strconv.Itoa(1), - "fileName": file.GetName(), - } - req.SetQueryParams(uploadParams) - req.SetBody(driver.NewLimitedUploadStream(ctx, rd)) - req.SetHeader("Content-Length", strconv.FormatInt(curChunkSize, 10)) - }, filePart); err != nil { - return err - } - - return nil - }, - retry.Attempts(3), - retry.DelayType(retry.BackOffDelay), - retry.Delay(time.Second)); err != nil { - return err - } - - if filePart.Name != "" { - fileParts = append(fileParts, *filePart) - uploaded += curChunkSize - up(float64(uploaded) / float64(totalSize)) - ss.FreeSectionReader(rd) - } - - } - - return d.createFileOnUploadSuccess(file.GetName(), fileId, dstDir.GetPath(), fileParts, totalSize) -} - -func (d *Teldrive) doMultiUpload(ctx context.Context, dstDir model.Obj, file model.FileStreamer, up model.UpdateProgress, - maxRetried, totalParts int, chunkSize int64, fileId string) error { - - concurrent := d.UploadConcurrency - g, ctx := errgroup.WithContext(ctx) - sem := semaphore.NewWeighted(int64(concurrent)) - chunkChan := make(chan chunkTask, concurrent*2) - resultChan := make(chan FilePart, concurrent) - totalSize := file.GetSize() - - ss, err := stream.NewStreamSectionReader(file, int(totalSize), &up) - if err != nil { - return err - } - ssLock := sync.Mutex{} - g.Go(func() error { - defer close(chunkChan) - - chunkIdx := 0 - for chunkIdx < totalParts { - select { - case <-ctx.Done(): - return ctx.Err() - default: - } - - offset := int64(chunkIdx) * chunkSize - curChunkSize := min(totalSize-offset, chunkSize) - - ssLock.Lock() - reader, err := ss.GetSectionReader(offset, curChunkSize) - ssLock.Unlock() - - if err != nil { - return err - } - task := chunkTask{ - chunkIdx: chunkIdx + 1, - chunkSize: curChunkSize, - fileName: file.GetName(), - reader: reader, - ss: ss, - } - // freeSectionReader will be called in d.uploadSingleChunk - select { - case chunkChan <- task: - chunkIdx++ - case <-ctx.Done(): - return ctx.Err() - } - } - return nil - }) - for i := 0; i < int(concurrent); i++ { - g.Go(func() error { - for task := range chunkChan { - if err := sem.Acquire(ctx, 1); err != nil { - return err - } - - filePart, err := d.uploadSingleChunk(ctx, fileId, task, totalParts, maxRetried) - sem.Release(1) - - if err != nil { - return fmt.Errorf("upload chunk %d failed: %w", task.chunkIdx, err) - } - - select { - case resultChan <- *filePart: - case <-ctx.Done(): - return ctx.Err() - } - } - return nil - }) - } - var fileParts []FilePart - var collectErr error - collectDone := make(chan struct{}) - - go func() { - defer close(collectDone) - fileParts = make([]FilePart, 0, totalParts) - - done := make(chan error, 1) - go func() { - done <- g.Wait() - close(resultChan) - }() - - for { - select { - case filePart, ok := <-resultChan: - if !ok { - collectErr = <-done - return - } - fileParts = append(fileParts, filePart) - case err := <-done: - collectErr = err - return - } - } - }() - - <-collectDone - - if collectErr != nil { - return fmt.Errorf("multi-upload failed: %w", collectErr) - } - sort.Slice(fileParts, func(i, j int) bool { - return fileParts[i].PartNo < fileParts[j].PartNo - }) - - return d.createFileOnUploadSuccess(file.GetName(), fileId, dstDir.GetPath(), fileParts, totalSize) -} - -func (d *Teldrive) uploadSingleChunk(ctx context.Context, fileId string, task chunkTask, totalParts, maxRetried int) (*FilePart, error) { - filePart := &FilePart{} - retryCount := 0 - defer task.ss.FreeSectionReader(task.reader) - - for { - select { - case <-ctx.Done(): - return nil, ctx.Err() - default: - } - - if existingPart, err := d.checkFilePartExist(fileId, task.chunkIdx); err == nil && existingPart.Name != "" { - return &existingPart, nil - } - - err := d.singleUploadRequest(fileId, func(req *resty.Request) { - uploadParams := map[string]string{ - "partName": func() string { - digits := len(fmt.Sprintf("%d", totalParts)) - return task.fileName + fmt.Sprintf(".%0*d", digits, task.chunkIdx) - }(), - "partNo": strconv.Itoa(task.chunkIdx), - "fileName": task.fileName, - } - req.SetQueryParams(uploadParams) - req.SetBody(driver.NewLimitedUploadStream(ctx, task.reader)) - req.SetHeader("Content-Length", strconv.Itoa(int(task.chunkSize))) - }, filePart) - - if err == nil { - return filePart, nil - } - - if retryCount >= maxRetried { - return nil, fmt.Errorf("upload failed after %d retries: %w", maxRetried, err) - } - - if errors.Is(err, context.DeadlineExceeded) || errors.Is(err, context.Canceled) { - continue - } - - retryCount++ - utils.Log.Errorf("[Teldrive] upload error: %v, retrying %d times", err, retryCount) - - backoffDuration := time.Duration(retryCount*retryCount) * time.Second - if backoffDuration > 30*time.Second { - backoffDuration = 30 * time.Second - } - - select { - case <-time.After(backoffDuration): - case <-ctx.Done(): - return nil, ctx.Err() - } - } -} - func (d *Teldrive) createShareFile(fileId string) error { var errResp ErrResp if err := d.request(http.MethodPost, "/api/files/"+fileId+"/share", func(req *resty.Request) { @@ -462,125 +104,3 @@ func getDateTime() string { formattedWithMs := now.Add(time.Hour * 1).Format("2006-01-02T15:04:05.000Z") return formattedWithMs } - -func NewCopyManager(ctx context.Context, concurrent int, d *Teldrive) *CopyManager { - g, ctx := errgroup.WithContext(ctx) - - return &CopyManager{ - TaskChan: make(chan CopyTask, concurrent*2), - Sem: semaphore.NewWeighted(int64(concurrent)), - G: g, - Ctx: ctx, - d: d, - } -} - -func (cm *CopyManager) startWorkers() { - workerCount := cap(cm.TaskChan) / 2 - for i := 0; i < workerCount; i++ { - cm.G.Go(func() error { - return cm.worker() - }) - } -} - -func (cm *CopyManager) worker() error { - for { - select { - case task, ok := <-cm.TaskChan: - if !ok { - return nil - } - - if err := cm.Sem.Acquire(cm.Ctx, 1); err != nil { - return err - } - - var err error - - err = cm.processFile(task) - - cm.Sem.Release(1) - - if err != nil { - return fmt.Errorf("task processing failed: %w", err) - } - - case <-cm.Ctx.Done(): - return cm.Ctx.Err() - } - } -} - -func (cm *CopyManager) generateTasks(ctx context.Context, srcObj, dstDir model.Obj) error { - if srcObj.IsDir() { - return cm.generateFolderTasks(ctx, srcObj, dstDir) - } else { - // add single file task directly - select { - case cm.TaskChan <- CopyTask{SrcObj: srcObj, DstDir: dstDir}: - return nil - case <-ctx.Done(): - return ctx.Err() - } - } -} - -func (cm *CopyManager) generateFolderTasks(ctx context.Context, srcDir, dstDir model.Obj) error { - objs, err := cm.d.List(ctx, srcDir, model.ListArgs{}) - if err != nil { - return fmt.Errorf("failed to list directory %s: %w", srcDir.GetPath(), err) - } - - err = cm.d.MakeDir(cm.Ctx, dstDir, srcDir.GetName()) - if err != nil || len(objs) == 0 { - return err - } - newDstDir := &model.Object{ - ID: dstDir.GetID(), - Path: dstDir.GetPath() + "/" + srcDir.GetName(), - Name: srcDir.GetName(), - IsFolder: true, - } - - for _, file := range objs { - if utils.IsCanceled(ctx) { - return ctx.Err() - } - - srcFile := &model.Object{ - ID: file.GetID(), - Path: srcDir.GetPath() + "/" + file.GetName(), - Name: file.GetName(), - IsFolder: file.IsDir(), - } - - // 递归生成任务 - if err := cm.generateTasks(ctx, srcFile, newDstDir); err != nil { - return err - } - } - - return nil -} - -func (cm *CopyManager) processFile(task CopyTask) error { - return cm.copySingleFile(cm.Ctx, task.SrcObj, task.DstDir) -} - -func (cm *CopyManager) copySingleFile(ctx context.Context, srcObj, dstDir model.Obj) error { - // `override copy mode` should delete the existing file - if obj, err := cm.d.getFile(dstDir.GetPath(), srcObj.GetName(), srcObj.IsDir()); err == nil { - if err := cm.d.Remove(ctx, obj); err != nil { - return fmt.Errorf("failed to remove existing file: %w", err) - } - } - - // Do copy - return cm.d.request(http.MethodPost, "/api/files/"+srcObj.GetID()+"/copy", func(req *resty.Request) { - req.SetBody(base.Json{ - "newName": srcObj.GetName(), - "destination": dstDir.GetPath(), - }) - }, nil) -} From 42256ee9025e8dba8f57053007ebc335a3196503 Mon Sep 17 00:00:00 2001 From: MadDogOwner Date: Sun, 24 Aug 2025 03:22:33 +0800 Subject: [PATCH 06/10] feat(teldrive): add UseShareLink option and support 302 --- drivers/teldrive/driver.go | 40 ++++++++++++++++---------------------- drivers/teldrive/meta.go | 4 ++-- 2 files changed, 19 insertions(+), 25 deletions(-) diff --git a/drivers/teldrive/driver.go b/drivers/teldrive/driver.go index 38d40bd69..b08924375 100644 --- a/drivers/teldrive/driver.go +++ b/drivers/teldrive/driver.go @@ -83,29 +83,23 @@ func (d *Teldrive) List(ctx context.Context, dir model.Obj, args model.ListArgs) } func (d *Teldrive) Link(ctx context.Context, file model.Obj, args model.LinkArgs) (*model.Link, error) { - //if d.WebdavPolicy != "native_proxy" { - // var address string - // if d.WebdavPolicy == "use_proxy_url" { - // address = d.DownProxyURL - // } else { - // address = d.Address - // } - // if shareObj, err := d.getShareFileById(file.GetID()); err == nil && shareObj != nil { - // return &model.Link{ - // URL: address + fmt.Sprintf("/api/shares/%s/files/%s/%s", shareObj.Id, file.GetID(), file.GetName()), - // }, nil - // } - // if err := d.createShareFile(file.GetID()); err != nil { - // return nil, err - // } - // shareObj, err := d.getShareFileById(file.GetID()) - // if err != nil { - // return nil, err - // } - // return &model.Link{ - // URL: address + fmt.Sprintf("/api/shares/%s/files/%s/%s", shareObj.Id, file.GetID(), file.GetName()), - // }, nil - //} + if d.UseShareLink { + if shareObj, err := d.getShareFileById(file.GetID()); err == nil && shareObj != nil { + return &model.Link{ + URL: d.Address + fmt.Sprintf("/api/shares/%s/files/%s/%s", shareObj.Id, file.GetID(), file.GetName()), + }, nil + } + if err := d.createShareFile(file.GetID()); err != nil { + return nil, err + } + shareObj, err := d.getShareFileById(file.GetID()) + if err != nil { + return nil, err + } + return &model.Link{ + URL: d.Address + fmt.Sprintf("/api/shares/%s/files/%s/%s", shareObj.Id, file.GetID(), file.GetName()), + }, nil + } return &model.Link{ URL: d.Address + "/api/files/" + file.GetID() + "/" + file.GetName(), Header: http.Header{ diff --git a/drivers/teldrive/meta.go b/drivers/teldrive/meta.go index cd25a7d18..028b6da84 100644 --- a/drivers/teldrive/meta.go +++ b/drivers/teldrive/meta.go @@ -10,15 +10,15 @@ type Addition struct { driver.RootPath // define other Address string `json:"url" required:"true"` - ChunkSize int64 `json:"chunk_size" type:"number" default:"4" help:"Chunk size in MiB"` Cookie string `json:"cookie" type:"string" required:"true" help:"access_token=xxx"` + UseShareLink bool `json:"use_share_link" type:"bool" default:"false" help:"Create share link when getting file link, support 302. If disabled, you need to enable web proxy."` + ChunkSize int64 `json:"chunk_size" type:"number" default:"4" help:"Chunk size in MiB"` UploadConcurrency int64 `json:"upload_concurrency" type:"number" default:"4" help:"Concurrency upload requests"` } var config = driver.Config{ Name: "Teldrive", DefaultRoot: "/", - OnlyProxy: true, } func init() { From 48ea347aff6805a4b591e548da549d8e4c733843 Mon Sep 17 00:00:00 2001 From: MadDogOwner Date: Sun, 24 Aug 2025 03:51:19 +0800 Subject: [PATCH 07/10] fix(teldrive): standardize API path construction --- drivers/teldrive/copy.go | 3 ++- drivers/teldrive/driver.go | 49 +++++++++++++++++++------------------- drivers/teldrive/meta.go | 2 -- drivers/teldrive/upload.go | 8 +++++-- drivers/teldrive/util.go | 7 ++++-- 5 files changed, 37 insertions(+), 32 deletions(-) diff --git a/drivers/teldrive/copy.go b/drivers/teldrive/copy.go index 6495d702d..1118f63a5 100644 --- a/drivers/teldrive/copy.go +++ b/drivers/teldrive/copy.go @@ -127,7 +127,8 @@ func (cm *CopyManager) copySingleFile(ctx context.Context, srcObj, dstDir model. } // Do copy - return cm.d.request(http.MethodPost, "/api/files/"+srcObj.GetID()+"/copy", func(req *resty.Request) { + return cm.d.request(http.MethodPost, "/api/files/{id}/copy", func(req *resty.Request) { + req.SetPathParam("id", srcObj.GetID()) req.SetBody(base.Json{ "newName": srcObj.GetName(), "destination": dstDir.GetPath(), diff --git a/drivers/teldrive/driver.go b/drivers/teldrive/driver.go index b08924375..8ad526959 100644 --- a/drivers/teldrive/driver.go +++ b/drivers/teldrive/driver.go @@ -51,17 +51,13 @@ func (d *Teldrive) Drop(ctx context.Context) error { } func (d *Teldrive) List(ctx context.Context, dir model.Obj, args model.ListArgs) ([]model.Obj, error) { - // endpoint /api/files, params ->page order sort path var listResp ListResp - params := url.Values{} - params.Set("path", dir.GetPath()) - //log.Info(dir.GetPath()) - pathname, err := utils.InjectQuery("/api/files", params) - if err != nil { - return nil, err - } - - err = d.request(http.MethodGet, pathname, nil, &listResp) + err := d.request(http.MethodGet, "/api/files", func(req *resty.Request) { + req.SetQueryParams(map[string]string{ + "path": dir.GetPath(), + "limit": "1000", // overide default 500, TODO pagination + }) + }, &listResp) if err != nil { return nil, err } @@ -84,24 +80,22 @@ func (d *Teldrive) List(ctx context.Context, dir model.Obj, args model.ListArgs) func (d *Teldrive) Link(ctx context.Context, file model.Obj, args model.LinkArgs) (*model.Link, error) { if d.UseShareLink { - if shareObj, err := d.getShareFileById(file.GetID()); err == nil && shareObj != nil { - return &model.Link{ - URL: d.Address + fmt.Sprintf("/api/shares/%s/files/%s/%s", shareObj.Id, file.GetID(), file.GetName()), - }, nil - } - if err := d.createShareFile(file.GetID()); err != nil { - return nil, err - } shareObj, err := d.getShareFileById(file.GetID()) - if err != nil { - return nil, err + if err != nil || shareObj == nil { + if err := d.createShareFile(file.GetID()); err != nil { + return nil, err + } + shareObj, err = d.getShareFileById(file.GetID()) + if err != nil { + return nil, err + } } return &model.Link{ - URL: d.Address + fmt.Sprintf("/api/shares/%s/files/%s/%s", shareObj.Id, file.GetID(), file.GetName()), + URL: d.Address + "/api/shares/" + url.PathEscape(shareObj.Id) + "/files/" + url.PathEscape(file.GetID()) + "/" + url.PathEscape(file.GetName()), }, nil } return &model.Link{ - URL: d.Address + "/api/files/" + file.GetID() + "/" + file.GetName(), + URL: d.Address + "/api/files/" + url.PathEscape(file.GetID()) + "/" + url.PathEscape(file.GetName()), Header: http.Header{ "Cookie": {d.Cookie}, }, @@ -130,7 +124,8 @@ func (d *Teldrive) Rename(ctx context.Context, srcObj model.Obj, newName string) body := base.Json{ "name": newName, } - return d.request(http.MethodPatch, "/api/files/"+srcObj.GetID(), func(req *resty.Request) { + return d.request(http.MethodPatch, "/api/files/{id}", func(req *resty.Request) { + req.SetPathParam("id", srcObj.GetID()) req.SetBody(body) }, nil) } @@ -165,7 +160,9 @@ func (d *Teldrive) Put(ctx context.Context, dstDir model.Obj, file model.FileStr // delete the upload task when finished or failed defer func() { - _ = d.request(http.MethodDelete, "/api/uploads/"+fileId, nil, nil) + _ = d.request(http.MethodDelete, "/api/uploads/{id}", func(req *resty.Request) { + req.SetPathParam("id", fileId) + }, nil) }() if obj, err := d.getFile(dstDir.GetPath(), file.GetName(), file.IsDir()); err == nil { @@ -174,7 +171,9 @@ func (d *Teldrive) Put(ctx context.Context, dstDir model.Obj, file model.FileStr } } // start the upload process - if err := d.request(http.MethodGet, "/api/uploads/"+fileId, nil, nil); err != nil { + if err := d.request(http.MethodGet, "/api/uploads/fileId", func(req *resty.Request) { + req.SetPathParam("id", fileId) + }, nil); err != nil { return err } if totalSize == 0 { diff --git a/drivers/teldrive/meta.go b/drivers/teldrive/meta.go index 028b6da84..19075f904 100644 --- a/drivers/teldrive/meta.go +++ b/drivers/teldrive/meta.go @@ -6,9 +6,7 @@ import ( ) type Addition struct { - // Usually one of two driver.RootPath - // define other Address string `json:"url" required:"true"` Cookie string `json:"cookie" type:"string" required:"true" help:"access_token=xxx"` UseShareLink bool `json:"use_share_link" type:"bool" default:"false" help:"Create share link when getting file link, support 302. If disabled, you need to enable web proxy."` diff --git a/drivers/teldrive/upload.go b/drivers/teldrive/upload.go index 7efb995b6..168d9beff 100644 --- a/drivers/teldrive/upload.go +++ b/drivers/teldrive/upload.go @@ -75,7 +75,9 @@ func (d *Teldrive) checkFilePartExist(fileId string, partId int) (FilePart, erro var uploadedParts []FilePart var filePart FilePart - if err := d.request(http.MethodGet, "/api/uploads/"+fileId, nil, &uploadedParts); err != nil { + if err := d.request(http.MethodGet, "/api/uploads/{id}", func(req *resty.Request) { + req.SetPathParam("id", fileId) + }, &uploadedParts); err != nil { return filePart, err } @@ -90,7 +92,9 @@ func (d *Teldrive) checkFilePartExist(fileId string, partId int) (FilePart, erro func (d *Teldrive) getFilePart(fileId string) ([]FilePart, error) { var uploadedParts []FilePart - if err := d.request(http.MethodGet, "/api/uploads/"+fileId, nil, &uploadedParts); err != nil { + if err := d.request(http.MethodGet, "/api/uploads/{id}", func(req *resty.Request) { + req.SetPathParam("id", fileId) + }, &uploadedParts); err != nil { return nil, err } diff --git a/drivers/teldrive/util.go b/drivers/teldrive/util.go index bda9d4506..ca3cccf97 100644 --- a/drivers/teldrive/util.go +++ b/drivers/teldrive/util.go @@ -75,7 +75,8 @@ func (err *ErrResp) Error() string { func (d *Teldrive) createShareFile(fileId string) error { var errResp ErrResp - if err := d.request(http.MethodPost, "/api/files/"+fileId+"/share", func(req *resty.Request) { + if err := d.request(http.MethodPost, "/api/files/{id}/share", func(req *resty.Request) { + req.SetPathParam("id", fileId) req.SetBody(base.Json{ "expiresAt": getDateTime(), }) @@ -92,7 +93,9 @@ func (d *Teldrive) createShareFile(fileId string) error { func (d *Teldrive) getShareFileById(fileId string) (*ShareObj, error) { var shareObj ShareObj - if err := d.request(http.MethodGet, "/api/files/"+fileId+"/share", nil, &shareObj); err != nil { + if err := d.request(http.MethodGet, "/api/files/{id}/share", func(req *resty.Request) { + req.SetPathParam("id", fileId) + }, &shareObj); err != nil { return nil, err } From 8f0e2b69739beb66ad8596b35b68a418fd8a0805 Mon Sep 17 00:00:00 2001 From: MadDogOwner Date: Sun, 24 Aug 2025 03:54:18 +0800 Subject: [PATCH 08/10] fix(teldrive): trim trailing slash from Address in Init method --- drivers/teldrive/driver.go | 1 + 1 file changed, 1 insertion(+) diff --git a/drivers/teldrive/driver.go b/drivers/teldrive/driver.go index 8ad526959..541d2e3be 100644 --- a/drivers/teldrive/driver.go +++ b/drivers/teldrive/driver.go @@ -32,6 +32,7 @@ func (d *Teldrive) GetAddition() driver.Additional { } func (d *Teldrive) Init(ctx context.Context) error { + d.Address = strings.TrimSuffix(d.Address, "/") if d.Cookie == "" || !strings.HasPrefix(d.Cookie, "access_token=") { return fmt.Errorf("cookie must start with 'access_token='") } From 7057a74c3fb7150a79cde877332f69c854f9b59f Mon Sep 17 00:00:00 2001 From: MadDogOwner Date: Sun, 24 Aug 2025 04:04:54 +0800 Subject: [PATCH 09/10] chore(teldrive): update help text for UseShareLink field in Addition struct --- drivers/teldrive/meta.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/drivers/teldrive/meta.go b/drivers/teldrive/meta.go index 19075f904..a16a78791 100644 --- a/drivers/teldrive/meta.go +++ b/drivers/teldrive/meta.go @@ -9,7 +9,7 @@ type Addition struct { driver.RootPath Address string `json:"url" required:"true"` Cookie string `json:"cookie" type:"string" required:"true" help:"access_token=xxx"` - UseShareLink bool `json:"use_share_link" type:"bool" default:"false" help:"Create share link when getting file link, support 302. If disabled, you need to enable web proxy."` + UseShareLink bool `json:"use_share_link" type:"bool" default:"false" help:"Create share link when getting link to support 302. If disabled, you need to enable web proxy."` ChunkSize int64 `json:"chunk_size" type:"number" default:"4" help:"Chunk size in MiB"` UploadConcurrency int64 `json:"upload_concurrency" type:"number" default:"4" help:"Concurrency upload requests"` } From 76b2bcb4adbea244123cf01b5dd1c1554aca22a9 Mon Sep 17 00:00:00 2001 From: twoonefour Date: Mon, 25 Aug 2025 00:52:47 +0800 Subject: [PATCH 10/10] fix(teldrive): set 10 MiB as default chunk size --- drivers/teldrive/meta.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/drivers/teldrive/meta.go b/drivers/teldrive/meta.go index a16a78791..23bae5f94 100644 --- a/drivers/teldrive/meta.go +++ b/drivers/teldrive/meta.go @@ -10,7 +10,7 @@ type Addition struct { Address string `json:"url" required:"true"` Cookie string `json:"cookie" type:"string" required:"true" help:"access_token=xxx"` UseShareLink bool `json:"use_share_link" type:"bool" default:"false" help:"Create share link when getting link to support 302. If disabled, you need to enable web proxy."` - ChunkSize int64 `json:"chunk_size" type:"number" default:"4" help:"Chunk size in MiB"` + ChunkSize int64 `json:"chunk_size" type:"number" default:"10" help:"Chunk size in MiB"` UploadConcurrency int64 `json:"upload_concurrency" type:"number" default:"4" help:"Concurrency upload requests"` }