mirror of
https://github.com/OpenListTeam/OpenList.git
synced 2026-10-10 21:13:10 +08:00
Compare commits
36 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 874234449b | |||
| 5fe267089a | |||
| 2442e302ad | |||
| 0612271732 | |||
| c261ce78fb | |||
| 7398e7d45e | |||
| 6e2d499ca9 | |||
| 4680ece2d9 | |||
| 8a4f3769d8 | |||
| cc5172e70b | |||
| a32ae97860 | |||
| d6dd62dfe5 | |||
| 216f071e64 | |||
| f47df5f9b2 | |||
| ff3c4b885c | |||
| f86c7c844c | |||
| 5db2172ed6 | |||
| c4c121befc | |||
| b4542753ba | |||
| 2a99c97d52 | |||
| 0a407c3d8b | |||
| b2596fdc24 | |||
| 2dbe1b00d3 | |||
| 1fc9c83df1 | |||
| d31e1a333d | |||
| e1bba7072b | |||
| 94c7d68413 | |||
| 9ed77a5875 | |||
| 7d6d3b8f55 | |||
| 5480d61f70 | |||
| 96cd714385 | |||
| e29d92f92e | |||
| c5f57bbcc5 | |||
| 9835afc645 | |||
| 1f373eac8d | |||
| ede96a314c |
@@ -122,12 +122,17 @@ Thank you for your support and understanding of the OpenList project.
|
||||
|
||||
## Demo
|
||||
|
||||
N/A (to be rebuilt)
|
||||
- 🌎 [Global Demo](https://demo.oplist.org)
|
||||
- 🇨🇳 [CN Demo](https://demo.oplist.org.cn)
|
||||
|
||||
## Discussion
|
||||
|
||||
Please refer to [*Discussions*](https://github.com/OpenListTeam/OpenList/discussions) for raising general questions, ***Issues* is for bug reports and feature requests only.**
|
||||
|
||||
## Sponsor
|
||||
|
||||
[](https://vps.town "VPS.Town - Trust, Effortlessly. Your Cloud, Reimagined.")
|
||||
|
||||
## License
|
||||
|
||||
The `OpenList` is open-source software licensed under the [AGPL-3.0](https://www.gnu.org/licenses/agpl-3.0.txt) license.
|
||||
|
||||
+6
-1
@@ -122,12 +122,17 @@ OpenList 是一个由 OpenList 团队独立维护的开源项目,遵循 AGPL-3
|
||||
|
||||
## 演示
|
||||
|
||||
N/A(待重建)
|
||||
- 🇨🇳 [国内演示站](https://demo.oplist.org.cn)
|
||||
- 🌎 [海外演示站](https://demo.oplist.org)
|
||||
|
||||
## 讨论
|
||||
|
||||
如有一般性问题请前往 [*Discussions*](https://github.com/OpenListTeam/OpenList/discussions) 讨论区,***Issues* 仅用于错误报告和功能请求。**
|
||||
|
||||
## 赞助者
|
||||
|
||||
[](https://vps.town "VPS.Town - Trust, Effortlessly. Your Cloud, Reimagined.")
|
||||
|
||||
## 许可证
|
||||
|
||||
`OpenList` 是基于 [AGPL-3.0](https://www.gnu.org/licenses/agpl-3.0.txt) 许可证的开源软件。
|
||||
|
||||
+6
-1
@@ -122,12 +122,17 @@ OpenListプロジェクトへのご支援とご理解をありがとうござい
|
||||
|
||||
## デモ
|
||||
|
||||
N/A(再構築中)
|
||||
- 🌎 [グローバルデモ](https://demo.oplist.org)
|
||||
- 🇨🇳 [CNデモ](https://demo.oplist.org.cn)
|
||||
|
||||
## ディスカッション
|
||||
|
||||
一般的な質問は [*Discussions*](https://github.com/OpenListTeam/OpenList/discussions) をご利用ください。***Issues* はバグ報告と機能リクエスト専用です。**
|
||||
|
||||
## スポンサー
|
||||
|
||||
[](https://vps.town "VPS.Town - Trust, Effortlessly. Your Cloud, Reimagined.")
|
||||
|
||||
## ライセンス
|
||||
|
||||
「OpenList」は [AGPL-3.0](https://www.gnu.org/licenses/agpl-3.0.txt) ライセンスの下で公開されているオープンソースソフトウェアです。
|
||||
|
||||
+6
-1
@@ -122,12 +122,17 @@ Dank u voor uw ondersteuning en begrip
|
||||
|
||||
## Demo
|
||||
|
||||
N.v.t. (wordt opnieuw opgebouwd)
|
||||
- 🌎 [Global Demo](https://demo.oplist.org)
|
||||
- 🇨🇳 [CN Demo](https://demo.oplist.org.cn)
|
||||
|
||||
## Discussie
|
||||
|
||||
Stel algemene vragen in [*Discussions*](https://github.com/OpenListTeam/OpenList/discussions), ***Issues* zijn alleen voor bugmeldingen en feature requests.**
|
||||
|
||||
## Sponsoren
|
||||
|
||||
[](https://vps.town "VPS.Town - Trust, Effortlessly. Your Cloud, Reimagined.")
|
||||
|
||||
## Licentie
|
||||
|
||||
`OpenList` is open-source software onder de [AGPL-3.0](https://www.gnu.org/licenses/agpl-3.0.txt) licentie.
|
||||
|
||||
+7
-6
@@ -6,6 +6,7 @@ package cmd
|
||||
import (
|
||||
"fmt"
|
||||
|
||||
"github.com/OpenListTeam/OpenList/v4/internal/bootstrap"
|
||||
"github.com/OpenListTeam/OpenList/v4/internal/conf"
|
||||
"github.com/OpenListTeam/OpenList/v4/internal/op"
|
||||
"github.com/OpenListTeam/OpenList/v4/internal/setting"
|
||||
@@ -20,8 +21,8 @@ var AdminCmd = &cobra.Command{
|
||||
Aliases: []string{"password"},
|
||||
Short: "Show admin user's info and some operations about admin user's password",
|
||||
Run: func(cmd *cobra.Command, args []string) {
|
||||
Init()
|
||||
defer Release()
|
||||
bootstrap.Init()
|
||||
defer bootstrap.Release()
|
||||
admin, err := op.GetAdmin()
|
||||
if err != nil {
|
||||
utils.Log.Errorf("failed get admin user: %+v", err)
|
||||
@@ -61,8 +62,8 @@ var ShowTokenCmd = &cobra.Command{
|
||||
Use: "token",
|
||||
Short: "Show admin token",
|
||||
Run: func(cmd *cobra.Command, args []string) {
|
||||
Init()
|
||||
defer Release()
|
||||
bootstrap.Init()
|
||||
defer bootstrap.Release()
|
||||
token := setting.GetStr(conf.Token)
|
||||
utils.Log.Infof("show admin token from CLI")
|
||||
fmt.Println("Admin token:", token)
|
||||
@@ -70,8 +71,8 @@ var ShowTokenCmd = &cobra.Command{
|
||||
}
|
||||
|
||||
func setAdminPassword(pwd string) {
|
||||
Init()
|
||||
defer Release()
|
||||
bootstrap.Init()
|
||||
defer bootstrap.Release()
|
||||
admin, err := op.GetAdmin()
|
||||
if err != nil {
|
||||
utils.Log.Errorf("failed get admin user: %+v", err)
|
||||
|
||||
+3
-2
@@ -6,6 +6,7 @@ package cmd
|
||||
import (
|
||||
"fmt"
|
||||
|
||||
"github.com/OpenListTeam/OpenList/v4/internal/bootstrap"
|
||||
"github.com/OpenListTeam/OpenList/v4/internal/op"
|
||||
"github.com/OpenListTeam/OpenList/v4/pkg/utils"
|
||||
"github.com/spf13/cobra"
|
||||
@@ -16,8 +17,8 @@ var Cancel2FACmd = &cobra.Command{
|
||||
Use: "cancel2fa",
|
||||
Short: "Delete 2FA of admin user",
|
||||
Run: func(cmd *cobra.Command, args []string) {
|
||||
Init()
|
||||
defer Release()
|
||||
bootstrap.Init()
|
||||
defer bootstrap.Release()
|
||||
admin, err := op.GetAdmin()
|
||||
if err != nil {
|
||||
utils.Log.Errorf("failed to get admin user: %+v", err)
|
||||
|
||||
+2
-10
@@ -6,24 +6,16 @@ import (
|
||||
"strconv"
|
||||
|
||||
"github.com/OpenListTeam/OpenList/v4/internal/bootstrap"
|
||||
"github.com/OpenListTeam/OpenList/v4/internal/bootstrap/data"
|
||||
"github.com/OpenListTeam/OpenList/v4/internal/db"
|
||||
"github.com/OpenListTeam/OpenList/v4/pkg/utils"
|
||||
log "github.com/sirupsen/logrus"
|
||||
)
|
||||
|
||||
func Init() {
|
||||
bootstrap.InitConfig()
|
||||
bootstrap.Log()
|
||||
bootstrap.InitDB()
|
||||
data.InitData()
|
||||
bootstrap.InitStreamLimit()
|
||||
bootstrap.InitIndex()
|
||||
bootstrap.InitUpgradePatch()
|
||||
bootstrap.Init()
|
||||
}
|
||||
|
||||
func Release() {
|
||||
db.Close()
|
||||
bootstrap.Release()
|
||||
}
|
||||
|
||||
var pid = -1
|
||||
|
||||
+2
-4
@@ -1,19 +1,17 @@
|
||||
package cmd
|
||||
|
||||
import (
|
||||
log "github.com/sirupsen/logrus"
|
||||
|
||||
"io"
|
||||
"os"
|
||||
"path"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
|
||||
"github.com/spf13/cobra"
|
||||
|
||||
rcCrypt "github.com/rclone/rclone/backend/crypt"
|
||||
"github.com/rclone/rclone/fs/config/configmap"
|
||||
"github.com/rclone/rclone/fs/config/obscure"
|
||||
log "github.com/sirupsen/logrus"
|
||||
"github.com/spf13/cobra"
|
||||
)
|
||||
|
||||
// encryption and decryption command format for Crypt driver
|
||||
|
||||
+20
-3
@@ -8,7 +8,6 @@ import (
|
||||
"fmt"
|
||||
"io"
|
||||
"os"
|
||||
"reflect"
|
||||
"strings"
|
||||
|
||||
_ "github.com/OpenListTeam/OpenList/v4/drivers"
|
||||
@@ -69,15 +68,33 @@ func writeFile(name string, data interface{}) {
|
||||
log.Errorf("failed to unmarshal json: %+v", err)
|
||||
return
|
||||
}
|
||||
if reflect.DeepEqual(oldData, newData) {
|
||||
if mergeJson(newData, oldData) {
|
||||
log.Infof("%s.json no changed, skip", name)
|
||||
} else {
|
||||
log.Infof("%s.json changed, update file", name)
|
||||
//log.Infof("old: %+v\nnew:%+v", oldData, data)
|
||||
utils.WriteJsonToFile(fmt.Sprintf("lang/%s.json", name), newData, true)
|
||||
utils.WriteJsonToFile(fmt.Sprintf("lang/%s.json", name), oldData, true)
|
||||
}
|
||||
}
|
||||
|
||||
func mergeJson(source, target map[string]interface{}) bool {
|
||||
equal := true
|
||||
for k, v := range source {
|
||||
tgtV, tgtOk := target[k]
|
||||
if !tgtOk {
|
||||
equal = false
|
||||
target[k] = v
|
||||
} else {
|
||||
srcMap, srcIsMap := v.(map[string]interface{})
|
||||
tgtMap, tgtIsMap := tgtV.(map[string]interface{})
|
||||
if srcIsMap && tgtIsMap {
|
||||
equal = mergeJson(srcMap, tgtMap) && equal
|
||||
}
|
||||
}
|
||||
}
|
||||
return equal
|
||||
}
|
||||
|
||||
func generateDriversJson() {
|
||||
drivers := make(Drivers)
|
||||
drivers["drivers"] = make(KV[interface{}])
|
||||
|
||||
+4
-239
File diff suppressed because it is too large
Load Diff
+7
-6
@@ -8,6 +8,7 @@ import (
|
||||
"os"
|
||||
"strconv"
|
||||
|
||||
"github.com/OpenListTeam/OpenList/v4/internal/bootstrap"
|
||||
"github.com/OpenListTeam/OpenList/v4/internal/db"
|
||||
"github.com/OpenListTeam/OpenList/v4/pkg/utils"
|
||||
"github.com/charmbracelet/bubbles/table"
|
||||
@@ -30,8 +31,8 @@ var disableStorageCmd = &cobra.Command{
|
||||
return fmt.Errorf("mount path is required")
|
||||
}
|
||||
mountPath := args[0]
|
||||
Init()
|
||||
defer Release()
|
||||
bootstrap.Init()
|
||||
defer bootstrap.Release()
|
||||
storage, err := db.GetStorageByMountPath(mountPath)
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to query storage: %+v", err)
|
||||
@@ -69,8 +70,8 @@ var deleteStorageCmd = &cobra.Command{
|
||||
}
|
||||
}
|
||||
|
||||
Init()
|
||||
defer Release()
|
||||
bootstrap.Init()
|
||||
defer bootstrap.Release()
|
||||
err = db.DeleteStorageById(uint(id))
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to delete storage by id: %+v", err)
|
||||
@@ -123,8 +124,8 @@ var listStorageCmd = &cobra.Command{
|
||||
Use: "list",
|
||||
Short: "List all storages",
|
||||
RunE: func(cmd *cobra.Command, args []string) error {
|
||||
Init()
|
||||
defer Release()
|
||||
bootstrap.Init()
|
||||
defer bootstrap.Release()
|
||||
storages, _, err := db.GetStorages(1, -1)
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to query storages: %+v", err)
|
||||
|
||||
@@ -53,6 +53,12 @@ func (d *Open115) Init(ctx context.Context) error {
|
||||
if d.Addition.LimitRate > 0 {
|
||||
d.limiter = rate.NewLimiter(rate.Limit(d.Addition.LimitRate), 1)
|
||||
}
|
||||
if d.PageSize <= 0 {
|
||||
d.PageSize = 200
|
||||
} else if d.PageSize > 1150 {
|
||||
d.PageSize = 1150
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -69,7 +75,7 @@ func (d *Open115) Drop(ctx context.Context) error {
|
||||
|
||||
func (d *Open115) List(ctx context.Context, dir model.Obj, args model.ListArgs) ([]model.Obj, error) {
|
||||
var res []model.Obj
|
||||
pageSize := int64(200)
|
||||
pageSize := int64(d.PageSize)
|
||||
offset := int64(0)
|
||||
for {
|
||||
if err := d.WaitLimit(ctx); err != nil {
|
||||
|
||||
@@ -12,6 +12,7 @@ type Addition struct {
|
||||
OrderBy string `json:"order_by" type:"select" options:"file_name,file_size,user_utime,file_type"`
|
||||
OrderDirection string `json:"order_direction" type:"select" options:"asc,desc"`
|
||||
LimitRate float64 `json:"limit_rate" type:"float" default:"1" help:"limit all api request rate ([limit]r/1s)"`
|
||||
PageSize int64 `json:"page_size" type:"number" default:"200" help:"list api per page size of 115open driver"`
|
||||
AccessToken string `json:"access_token" required:"true"`
|
||||
RefreshToken string `json:"refresh_token" required:"true"`
|
||||
}
|
||||
|
||||
@@ -107,16 +107,16 @@ func (d *Open115) multpartUpload(ctx context.Context, stream model.FileStreamer,
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
rateLimitedRd := driver.NewLimitedUploadStream(ctx, rd)
|
||||
err = retry.Do(func() error {
|
||||
rd.Seek(0, io.SeekStart)
|
||||
part, err := bucket.UploadPart(imur, rateLimitedRd, partSize, int(i))
|
||||
part, err := bucket.UploadPart(imur, driver.NewLimitedUploadStream(ctx, rd), partSize, int(i))
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
parts[i-1] = part
|
||||
return nil
|
||||
},
|
||||
retry.Context(ctx),
|
||||
retry.Attempts(3),
|
||||
retry.DelayType(retry.BackOffDelay),
|
||||
retry.Delay(time.Second))
|
||||
|
||||
+7
-16
@@ -125,27 +125,18 @@ func (d *Pan123) newUpload(ctx context.Context, upReq *UploadResp, file model.Fi
|
||||
curSize = lastChunkSize
|
||||
}
|
||||
var reader io.ReadSeeker
|
||||
var rateLimitedRd io.Reader
|
||||
threadG.GoWithLifecycle(errgroup.Lifecycle{
|
||||
Before: func(ctx context.Context) error {
|
||||
if reader == nil {
|
||||
var err error
|
||||
reader, err = ss.GetSectionReader(offset, curSize)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
rateLimitedRd = driver.NewLimitedUploadStream(ctx, reader)
|
||||
}
|
||||
return nil
|
||||
Before: func(ctx context.Context) (err error) {
|
||||
reader, err = ss.GetSectionReader(offset, curSize)
|
||||
return
|
||||
},
|
||||
Do: func(ctx context.Context) error {
|
||||
Do: func(ctx context.Context) (err error) {
|
||||
reader.Seek(0, io.SeekStart)
|
||||
uploadUrl := s3PreSignedUrls.Data.PreSignedUrls[strconv.Itoa(cur)]
|
||||
if uploadUrl == "" {
|
||||
return fmt.Errorf("upload url is empty, s3PreSignedUrls: %+v", s3PreSignedUrls)
|
||||
}
|
||||
reader.Seek(0, io.SeekStart)
|
||||
req, err := http.NewRequestWithContext(ctx, http.MethodPut, uploadUrl, rateLimitedRd)
|
||||
req, err := http.NewRequestWithContext(ctx, http.MethodPut, uploadUrl, driver.NewLimitedUploadStream(ctx, reader))
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
@@ -157,7 +148,7 @@ func (d *Pan123) newUpload(ctx context.Context, upReq *UploadResp, file model.Fi
|
||||
}
|
||||
defer res.Body.Close()
|
||||
if res.StatusCode == http.StatusForbidden {
|
||||
singleflight.AnyGroup.Do(fmt.Sprintf("Pan123.newUpload_%p", threadG), func() (any, error) {
|
||||
_, err, _ = singleflight.AnyGroup.Do(fmt.Sprintf("Pan123.newUpload_%p", threadG), func() (any, error) {
|
||||
newS3PreSignedUrls, err := getS3UploadUrl(ctx, upReq, cur, end)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
@@ -177,7 +168,7 @@ func (d *Pan123) newUpload(ctx context.Context, upReq *UploadResp, file model.Fi
|
||||
}
|
||||
return fmt.Errorf("upload s3 chunk %d failed, status code: %d, body: %s", cur, res.StatusCode, body)
|
||||
}
|
||||
progress := 10.0 + 85.0*float64(threadG.Success())/float64(chunkCount)
|
||||
progress := 100 * float64(threadG.Success()+1) / float64(chunkCount+1)
|
||||
up(progress)
|
||||
return nil
|
||||
},
|
||||
|
||||
@@ -39,6 +39,10 @@ func (d *Pan123Link) Drop(ctx context.Context) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
func (Addition) GetRootPath() string {
|
||||
return "/"
|
||||
}
|
||||
|
||||
func (d *Pan123Link) Get(ctx context.Context, path string) (model.Obj, error) {
|
||||
node := GetNodeFromRootByPath(d.root, path)
|
||||
return nodeToObj(node, path)
|
||||
|
||||
@@ -18,6 +18,7 @@ type Open123 struct {
|
||||
model.Storage
|
||||
Addition
|
||||
UID uint64
|
||||
tm *tokenManager
|
||||
}
|
||||
|
||||
func (d *Open123) Config() driver.Config {
|
||||
@@ -33,6 +34,24 @@ func (d *Open123) Init(ctx context.Context) error {
|
||||
d.UploadThread = 3
|
||||
}
|
||||
|
||||
if d.RefreshToken != "" {
|
||||
// refresh token 直接主动刷新
|
||||
d.AccessToken = ""
|
||||
d.tm = &tokenManager{}
|
||||
} else {
|
||||
// 避免个人 token 刷新产生的多个登录,被动刷新
|
||||
// 默认过期时间90天,jwt exp 不可靠
|
||||
d.tm = &tokenManager{
|
||||
// accessToken: d.AccessToken,
|
||||
expiredAt: time.Now().Add(90 * 24 * time.Hour),
|
||||
}
|
||||
}
|
||||
|
||||
_, err := d.getAccessToken(false)
|
||||
if err != nil {
|
||||
return fmt.Errorf("init get access token error: %w", err)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
|
||||
@@ -13,7 +13,7 @@ type Addition struct {
|
||||
ClientID string `json:"ClientID" required:"false"`
|
||||
ClientSecret string `json:"ClientSecret" required:"false"`
|
||||
|
||||
// 直接写入AccessToken
|
||||
// 直接写入AccessToken, AccessToken有过期时间,不建议直接填写
|
||||
AccessToken string `json:"AccessToken" required:"false"`
|
||||
|
||||
// 用户名+密码方式登录的AccessToken可以兼容
|
||||
|
||||
@@ -0,0 +1,115 @@
|
||||
package _123_open
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"net/http"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/OpenListTeam/OpenList/v4/drivers/base"
|
||||
"github.com/OpenListTeam/OpenList/v4/internal/op"
|
||||
)
|
||||
|
||||
var (
|
||||
AccessToken = "https://open-api.123pan.com/api/v1/access_token"
|
||||
RefreshToken = "https://open-api.123pan.com/api/v1/oauth2/access_token"
|
||||
)
|
||||
|
||||
type tokenManager struct {
|
||||
// accessToken string
|
||||
expiredAt time.Time
|
||||
mu sync.Mutex
|
||||
blockRefresh bool
|
||||
}
|
||||
|
||||
func (d *Open123) getAccessToken(forceRefresh bool) (string, error) {
|
||||
tm := d.tm
|
||||
tm.mu.Lock()
|
||||
defer tm.mu.Unlock()
|
||||
if tm.blockRefresh {
|
||||
return "", errors.New("Authentication expired")
|
||||
}
|
||||
if !forceRefresh && d.AccessToken != "" && time.Now().Before(tm.expiredAt.Add(-5*time.Minute)) {
|
||||
return d.AccessToken, nil
|
||||
}
|
||||
if err := d.flushAccessToken(); err != nil {
|
||||
// token expired and failed to refresh, block further refresh attempts
|
||||
tm.blockRefresh = true
|
||||
return "", err
|
||||
}
|
||||
return d.AccessToken, nil
|
||||
}
|
||||
|
||||
func (d *Open123) flushAccessToken() error {
|
||||
// directly send request to avoid deadlock
|
||||
req := base.RestyClient.R()
|
||||
req.SetHeaders(map[string]string{
|
||||
"authorization": "Bearer " + d.AccessToken,
|
||||
"platform": "open_platform",
|
||||
"Content-Type": "application/json",
|
||||
})
|
||||
|
||||
if d.ClientID != "" {
|
||||
if d.RefreshToken != "" {
|
||||
var resp RefreshTokenResp
|
||||
req.SetQueryParam("client_id", d.ClientID)
|
||||
if d.ClientSecret != "" {
|
||||
req.SetQueryParam("client_secret", d.ClientSecret)
|
||||
}
|
||||
req.SetQueryParam("grant_type", "refresh_token")
|
||||
req.SetQueryParam("refresh_token", d.RefreshToken)
|
||||
req.SetResult(&resp)
|
||||
res, err := req.Execute(http.MethodPost, RefreshToken)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
body := res.Body()
|
||||
var baseResp BaseResp
|
||||
if err = json.Unmarshal(body, &baseResp); err != nil {
|
||||
return err
|
||||
}
|
||||
if baseResp.Code != 0 {
|
||||
return fmt.Errorf("get access token failed: %s", baseResp.Message)
|
||||
}
|
||||
|
||||
d.AccessToken = resp.AccessToken
|
||||
// add token expire time
|
||||
d.tm.expiredAt = time.Now().Add(time.Duration(resp.ExpiresIn) * time.Second)
|
||||
d.RefreshToken = resp.RefreshToken
|
||||
op.MustSaveDriverStorage(d)
|
||||
d.tm.blockRefresh = false
|
||||
return nil
|
||||
} else if d.ClientSecret != "" {
|
||||
var resp AccessTokenResp
|
||||
req.SetBody(base.Json{
|
||||
"clientID": d.ClientID,
|
||||
"clientSecret": d.ClientSecret,
|
||||
})
|
||||
req.SetResult(&resp)
|
||||
res, err := req.Execute(http.MethodPost, AccessToken)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
body := res.Body()
|
||||
var baseResp BaseResp
|
||||
if err = json.Unmarshal(body, &baseResp); err != nil {
|
||||
return err
|
||||
}
|
||||
if baseResp.Code != 0 {
|
||||
return fmt.Errorf("get access token failed: %s", baseResp.Message)
|
||||
}
|
||||
d.AccessToken = resp.Data.AccessToken
|
||||
// parse token expire time
|
||||
d.tm.expiredAt, err = time.Parse(time.RFC3339, resp.Data.ExpiredAt)
|
||||
if err != nil {
|
||||
return fmt.Errorf("parse expire time failed: %w", err)
|
||||
}
|
||||
op.MustSaveDriverStorage(d)
|
||||
d.tm.blockRefresh = false
|
||||
return nil
|
||||
}
|
||||
}
|
||||
return errors.New("no valid authentication method available")
|
||||
}
|
||||
+19
-19
@@ -73,25 +73,20 @@ func (d *Open123) Upload(ctx context.Context, file model.FileStreamer, createRes
|
||||
// 表单
|
||||
b := bytes.NewBuffer(make([]byte, 0, 2048))
|
||||
threadG.GoWithLifecycle(errgroup.Lifecycle{
|
||||
Before: func(ctx context.Context) error {
|
||||
if reader == nil {
|
||||
var err error
|
||||
// 每个分片一个reader
|
||||
reader, err = ss.GetSectionReader(offset, size)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
// 计算当前分片的MD5
|
||||
Before: func(ctx context.Context) (err error) {
|
||||
reader, err = ss.GetSectionReader(offset, size)
|
||||
return
|
||||
},
|
||||
Do: func(ctx context.Context) (err error) {
|
||||
reader.Seek(0, io.SeekStart)
|
||||
if sliceMD5 == "" {
|
||||
// 把耗时的计算放在这里,避免阻塞其他协程
|
||||
sliceMD5, err = utils.HashReader(utils.MD5, reader)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
reader.Seek(0, io.SeekStart)
|
||||
}
|
||||
return nil
|
||||
},
|
||||
Do: func(ctx context.Context) error {
|
||||
// 重置分片reader位置,因为HashReader、上一次失败已经读取到分片EOF
|
||||
reader.Seek(0, io.SeekStart)
|
||||
|
||||
b.Reset()
|
||||
w := multipart.NewWriter(b)
|
||||
@@ -121,6 +116,10 @@ func (d *Open123) Upload(ctx context.Context, file model.FileStreamer, createRes
|
||||
head := bytes.NewReader(b.Bytes()[:headSize])
|
||||
tail := bytes.NewReader(b.Bytes()[headSize:])
|
||||
rateLimitedRd = driver.NewLimitedUploadStream(ctx, io.MultiReader(head, reader, tail))
|
||||
token, err := d.getAccessToken(false)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
// 创建请求并设置header
|
||||
req, err := http.NewRequestWithContext(ctx, http.MethodPost, uploadDomain+"/upload/v2/file/slice", rateLimitedRd)
|
||||
if err != nil {
|
||||
@@ -128,7 +127,7 @@ func (d *Open123) Upload(ctx context.Context, file model.FileStreamer, createRes
|
||||
}
|
||||
|
||||
// 设置请求头
|
||||
req.Header.Add("Authorization", "Bearer "+d.AccessToken)
|
||||
req.Header.Add("Authorization", "Bearer "+token)
|
||||
req.Header.Add("Content-Type", w.FormDataContentType())
|
||||
req.Header.Add("Platform", "open_platform")
|
||||
|
||||
@@ -140,12 +139,13 @@ func (d *Open123) Upload(ctx context.Context, file model.FileStreamer, createRes
|
||||
if res.StatusCode != 200 {
|
||||
return fmt.Errorf("slice %d upload failed, status code: %d", partNumber, res.StatusCode)
|
||||
}
|
||||
var resp BaseResp
|
||||
respBody, err := io.ReadAll(res.Body)
|
||||
b.Reset()
|
||||
_, err = b.ReadFrom(res.Body)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
err = json.Unmarshal(respBody, &resp)
|
||||
var resp BaseResp
|
||||
err = json.Unmarshal(b.Bytes(), &resp)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
@@ -153,7 +153,7 @@ func (d *Open123) Upload(ctx context.Context, file model.FileStreamer, createRes
|
||||
return fmt.Errorf("slice %d upload failed: %s", partNumber, resp.Message)
|
||||
}
|
||||
|
||||
progress := 10.0 + 85.0*float64(threadG.Success())/float64(uploadNums)
|
||||
progress := 100 * float64(threadG.Success()+1) / float64(uploadNums+1)
|
||||
up(progress)
|
||||
return nil
|
||||
},
|
||||
|
||||
Some files were not shown because too many files have changed in this diff Show More
Reference in New Issue
Block a user