mirror of
https://github.com/OpenListTeam/OpenList.git
synced 2026-10-10 04:53:09 +08:00
4c39bbe9c2
QR code scans were authorized on the phone but the storage stayed stuck on the QR page, and token-only setups failed with "params is null". The driver used appId 8025431004 while the official PC client uses 9317140619. Tokens, sessions and QR sessions are all scoped to an appId, so the mismatch meant the QR poll never saw status:0 and an accessToken could not be exchanged for a sessionSecret. - Use appId 9317140619 and version 7.2.4.0. QR state polling uses clientType=1; password login keeps 10020. - Send the QR poll parameters the official client sends (cb_SaveName, isOauth2, state, user-finger header, logbox Referer) and poll locally instead of only checking once per save. - Parse lt/reqId from the logbox redirect and paramId from appConf.do. The new login page no longer embeds them as inline variables; the old inline format is still supported. - Implement the -133 second device verification via sendSmsCodeForSecondAuth/submitForSecondAuth. That endpoint has no dedicated SMS field: the code goes into epd, encrypted with the same public key used for the password. Persist the DEVICEID cookie so the verification only happens once. - Detect refreshToken.do failures. It reports them as HTTP 200 with a result field, so SetError never fired and a failed refresh was treated as success, surfacing later as a misleading "params is null". - username/password are no longer required, so token-only storages save without placeholders. clientSn/jgOpenId are optional and only sent when configured; user-finger is generated once per storage. - Return named errors instead of panicking when the login page changes shape. Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
1899 lines
55 KiB
Go
1899 lines
55 KiB
Go
package _189pc
|
||
|
||
import (
|
||
"bytes"
|
||
"context"
|
||
sha1Pkg "crypto/sha1"
|
||
"encoding/base64"
|
||
"encoding/hex"
|
||
"encoding/xml"
|
||
"fmt"
|
||
"hash"
|
||
"io"
|
||
"net/http"
|
||
"net/http/cookiejar"
|
||
"net/url"
|
||
"os"
|
||
"regexp"
|
||
"sort"
|
||
"strconv"
|
||
"strings"
|
||
"time"
|
||
|
||
"github.com/OpenListTeam/OpenList/v4/drivers/base"
|
||
"github.com/OpenListTeam/OpenList/v4/internal/conf"
|
||
"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/internal/setting"
|
||
"github.com/OpenListTeam/OpenList/v4/internal/stream"
|
||
"github.com/OpenListTeam/OpenList/v4/pkg/errgroup"
|
||
"github.com/OpenListTeam/OpenList/v4/pkg/utils"
|
||
"github.com/OpenListTeam/OpenList/v4/pkg/utils/random"
|
||
"github.com/skip2/go-qrcode"
|
||
|
||
"github.com/avast/retry-go"
|
||
"github.com/go-resty/resty/v2"
|
||
"github.com/google/uuid"
|
||
jsoniter "github.com/json-iterator/go"
|
||
"github.com/pkg/errors"
|
||
)
|
||
|
||
const (
|
||
ACCOUNT_TYPE = "02"
|
||
// 官方 PC 端(cloud.189.cn 网页/客户端)使用的 appId,
|
||
// 登录、生成二维码、换取 session 必须全程使用同一个 appId
|
||
APP_ID = "9317140619"
|
||
CLIENT_TYPE = "10020"
|
||
// 扫码状态轮询使用的 clientType,与密码登录的 10020 不同
|
||
QR_CLIENT_TYPE = "1"
|
||
VERSION = "7.2.4.0"
|
||
|
||
WEB_URL = "https://cloud.189.cn"
|
||
AUTH_URL = "https://open.e.189.cn"
|
||
API_URL = "https://api.cloud.189.cn"
|
||
UPLOAD_URL = "https://upload.cloud.189.cn"
|
||
|
||
RETURN_URL = "https://m.cloud.189.cn/zhuanti/2020/loginErrorPc/index.html"
|
||
|
||
PC = "TELEPC"
|
||
MAC = "TELEMAC"
|
||
|
||
CHANNEL_ID = "web_cloud.189.cn"
|
||
|
||
// 服务端通过短信二次校验后下发的设备标识,复用它可以避免再次触发校验
|
||
DEVICE_ID_COOKIE = "DEVICEID"
|
||
|
||
// 扫码登录本地轮询参数,超时后把二维码交回前端,避免请求被反向代理掐断
|
||
QRCODE_POLL_INTERVAL = 2 * time.Second
|
||
QRCODE_POLL_TIMEOUT = 20 * time.Second
|
||
|
||
// Error codes
|
||
UserInvalidOpenTokenError = "UserInvalidOpenToken"
|
||
|
||
// 密码登录返回该结果表示需要设备二次校验
|
||
SecondDeviceAuthResult = -133
|
||
)
|
||
|
||
func (y *Cloud189PC) SignatureHeader(url, method, params string, isFamily bool) map[string]string {
|
||
dateOfGmt := getHttpDateStr()
|
||
sessionKey := y.getTokenInfo().SessionKey
|
||
sessionSecret := y.getTokenInfo().SessionSecret
|
||
if isFamily {
|
||
sessionKey = y.getTokenInfo().FamilySessionKey
|
||
sessionSecret = y.getTokenInfo().FamilySessionSecret
|
||
}
|
||
|
||
header := map[string]string{
|
||
"Date": dateOfGmt,
|
||
"SessionKey": sessionKey,
|
||
"X-Request-ID": uuid.NewString(),
|
||
"Signature": signatureOfHmac(sessionSecret, sessionKey, method, url, dateOfGmt, params),
|
||
}
|
||
return header
|
||
}
|
||
|
||
func (y *Cloud189PC) EncryptParams(params Params, isFamily bool) string {
|
||
sessionSecret := y.getTokenInfo().SessionSecret
|
||
if isFamily {
|
||
sessionSecret = y.getTokenInfo().FamilySessionSecret
|
||
}
|
||
if params != nil {
|
||
return AesECBEncrypt(params.Encode(), sessionSecret[:16])
|
||
}
|
||
return ""
|
||
}
|
||
|
||
func (y *Cloud189PC) request(url, method string, callback base.ReqCallback, params Params, resp interface{}, isFamily ...bool) ([]byte, error) {
|
||
if y.getTokenInfo() == nil {
|
||
return nil, fmt.Errorf("login failed")
|
||
}
|
||
req := y.getClient().R().SetQueryParams(clientSuffix())
|
||
|
||
// 设置params
|
||
paramsData := y.EncryptParams(params, isBool(isFamily...))
|
||
if paramsData != "" {
|
||
req.SetQueryParam("params", paramsData)
|
||
}
|
||
|
||
// Signature
|
||
req.SetHeaders(y.SignatureHeader(url, method, paramsData, isBool(isFamily...)))
|
||
|
||
var erron RespErr
|
||
req.SetError(&erron)
|
||
|
||
if callback != nil {
|
||
callback(req)
|
||
}
|
||
if resp != nil {
|
||
req.SetResult(resp)
|
||
}
|
||
res, err := req.Execute(method, url)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
|
||
if strings.Contains(res.String(), "userSessionBO is null") {
|
||
if err = y.refreshSession(); err != nil {
|
||
return nil, err
|
||
}
|
||
return y.request(url, method, callback, params, resp, isFamily...)
|
||
}
|
||
|
||
// if erron.ErrorCode == "InvalidSessionKey" || erron.Code == "InvalidSessionKey" {
|
||
if strings.Contains(res.String(), "InvalidSessionKey") {
|
||
if err = y.refreshSession(); err != nil {
|
||
return nil, err
|
||
}
|
||
return y.request(url, method, callback, params, resp, isFamily...)
|
||
}
|
||
|
||
// 处理错误
|
||
if erron.HasError() {
|
||
return nil, &erron
|
||
}
|
||
return res.Body(), nil
|
||
}
|
||
|
||
func (y *Cloud189PC) get(url string, callback base.ReqCallback, resp interface{}, isFamily ...bool) ([]byte, error) {
|
||
return y.request(url, http.MethodGet, callback, nil, resp, isFamily...)
|
||
}
|
||
|
||
func (y *Cloud189PC) post(url string, callback base.ReqCallback, resp interface{}, isFamily ...bool) ([]byte, error) {
|
||
return y.request(url, http.MethodPost, callback, nil, resp, isFamily...)
|
||
}
|
||
|
||
func (y *Cloud189PC) put(ctx context.Context, url string, headers map[string]string, sign bool, file io.Reader, isFamily bool) ([]byte, error) {
|
||
req, err := http.NewRequestWithContext(ctx, http.MethodPut, url, file)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
|
||
query := req.URL.Query()
|
||
for key, value := range clientSuffix() {
|
||
query.Add(key, value)
|
||
}
|
||
req.URL.RawQuery = query.Encode()
|
||
|
||
for key, value := range headers {
|
||
req.Header.Add(key, value)
|
||
}
|
||
|
||
if sign {
|
||
for key, value := range y.SignatureHeader(url, http.MethodPut, "", isFamily) {
|
||
req.Header.Add(key, value)
|
||
}
|
||
}
|
||
|
||
resp, err := base.HttpClient.Do(req)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
defer resp.Body.Close()
|
||
|
||
body, err := io.ReadAll(resp.Body)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
|
||
var erron RespErr
|
||
_ = jsoniter.Unmarshal(body, &erron)
|
||
_ = xml.Unmarshal(body, &erron)
|
||
if erron.HasError() {
|
||
return nil, &erron
|
||
}
|
||
if resp.StatusCode != http.StatusOK {
|
||
return nil, errors.Errorf("put fail,err:%s", string(body))
|
||
}
|
||
return body, nil
|
||
}
|
||
|
||
func (y *Cloud189PC) getFiles(ctx context.Context, fileId string, isFamily bool) ([]model.Obj, error) {
|
||
pageSize := 1000 // 每一页返回的文件数量
|
||
res := make([]model.Obj, 0, 100)
|
||
for pageNum := 1; ; pageNum++ {
|
||
resp, err := y.getFilesWithPage(ctx, fileId, isFamily, pageNum, pageSize, y.OrderBy, y.OrderDirection)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
// 获取完毕跳出
|
||
if resp.FileListAO.Count == 0 {
|
||
break
|
||
}
|
||
|
||
FolderCount := len(resp.FileListAO.FolderList) // 当前文件夹总数
|
||
FileCount := len(resp.FileListAO.FileList) // 当前文件总数
|
||
PageCount := FolderCount + FileCount // 当前页数总数
|
||
|
||
for i := 0; i < FolderCount; i++ {
|
||
res = append(res, &resp.FileListAO.FolderList[i])
|
||
}
|
||
for i := 0; i < FileCount; i++ {
|
||
resp.FileListAO.FileList[i].ParentID = fileId
|
||
res = append(res, &resp.FileListAO.FileList[i])
|
||
}
|
||
|
||
// 当前文件数量小于设定数量则跳出
|
||
if PageCount < pageSize {
|
||
break
|
||
}
|
||
}
|
||
return res, nil
|
||
}
|
||
|
||
func (y *Cloud189PC) getFilesWithPage(ctx context.Context, fileId string, isFamily bool, pageNum int, pageSize int, orderBy string, orderDirection string) (*Cloud189FilesResp, error) {
|
||
fullUrl := API_URL
|
||
if isFamily {
|
||
fullUrl += "/family/file"
|
||
}
|
||
fullUrl += "/listFiles.action"
|
||
|
||
var resp Cloud189FilesResp
|
||
_, err := y.get(fullUrl, func(r *resty.Request) {
|
||
r.SetContext(ctx)
|
||
r.SetQueryParams(map[string]string{
|
||
"folderId": fileId,
|
||
"fileType": "0",
|
||
"mediaAttr": "0",
|
||
"iconOption": "5",
|
||
"pageNum": fmt.Sprint(pageNum),
|
||
"pageSize": fmt.Sprint(pageSize),
|
||
})
|
||
if isFamily {
|
||
r.SetQueryParams(map[string]string{
|
||
"familyId": y.FamilyID,
|
||
"orderBy": toFamilyOrderBy(orderBy),
|
||
"descending": toDesc(orderDirection),
|
||
})
|
||
} else {
|
||
r.SetQueryParams(map[string]string{
|
||
"recursive": "0",
|
||
"orderBy": orderBy,
|
||
"descending": toDesc(orderDirection),
|
||
})
|
||
}
|
||
}, &resp, isFamily)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
return &resp, nil
|
||
}
|
||
|
||
func (y *Cloud189PC) findFileByName(ctx context.Context, searchName string, folderId string, isFamily bool) (*Cloud189File, error) {
|
||
for pageNum := 1; ; pageNum++ {
|
||
resp, err := y.getFilesWithPage(ctx, folderId, isFamily, pageNum, 10, "filename", "asc")
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
// 获取完毕跳出
|
||
if resp.FileListAO.Count == 0 {
|
||
return nil, errs.ObjectNotFound
|
||
}
|
||
for i := 0; i < len(resp.FileListAO.FileList); i++ {
|
||
file := resp.FileListAO.FileList[i]
|
||
if file.Name == searchName {
|
||
return &file, nil
|
||
}
|
||
}
|
||
}
|
||
}
|
||
|
||
func (y *Cloud189PC) login() error {
|
||
if y.LoginType == "qrcode" {
|
||
return y.loginByQRCode()
|
||
}
|
||
if y.Username == "" || y.Password == "" {
|
||
return errors.New("please fill in the username and password, or provide an access token / refresh token")
|
||
}
|
||
return y.loginByPassword()
|
||
}
|
||
|
||
// 设备指纹,为空时生成并保存,服务端以此识别是否为同一台设备
|
||
func (y *Cloud189PC) getUserFinger() string {
|
||
if y.Addition.UserFinger == "" {
|
||
y.Addition.UserFinger = fmt.Sprint(random.Rand.Int63n(9e9) + 1e9)
|
||
op.MustSaveDriverStorage(y)
|
||
}
|
||
return y.Addition.UserFinger
|
||
}
|
||
|
||
// 换取会话时携带的设备参数,与官方PC客户端保持一致
|
||
// clientSn/jgOpenId 只在用户从官方客户端抓到并填写后才发送,避免上报一个服务端不认识的设备号
|
||
func (y *Cloud189PC) deviceParams() map[string]string {
|
||
params := map[string]string{"returnType": "JSON"}
|
||
if y.Addition.ClientSn != "" {
|
||
params["clientSn"] = y.Addition.ClientSn
|
||
}
|
||
if y.Addition.JgOpenId != "" {
|
||
params["jgOpenId"] = y.Addition.JgOpenId
|
||
}
|
||
return params
|
||
}
|
||
|
||
// logbox接口的公共请求头,缺少user-finger和Referer会被判定为陌生设备
|
||
func (y *Cloud189PC) loginHeaders(param BaseLoginParam) map[string]string {
|
||
return map[string]string{
|
||
"REQID": param.ReqId,
|
||
"lt": param.Lt,
|
||
"user-finger": y.getUserFinger(),
|
||
"Referer": IF(param.Referer != "", param.Referer, AUTH_URL),
|
||
}
|
||
}
|
||
|
||
// 把已保存的设备标识写入cookie,避免重复触发设备二次校验
|
||
func (y *Cloud189PC) applyDeviceID(jar http.CookieJar) {
|
||
if y.Addition.DeviceID == "" {
|
||
return
|
||
}
|
||
authUrl, err := url.Parse(AUTH_URL)
|
||
if err != nil {
|
||
return
|
||
}
|
||
jar.SetCookies(authUrl, []*http.Cookie{{
|
||
Name: DEVICE_ID_COOKIE,
|
||
Value: y.Addition.DeviceID,
|
||
Domain: "e.189.cn",
|
||
Path: "/",
|
||
}})
|
||
}
|
||
|
||
// 保存服务端下发的设备标识,下次登陆复用即可跳过设备二次校验
|
||
func (y *Cloud189PC) saveDeviceID(res *resty.Response) {
|
||
for _, cookie := range res.Cookies() {
|
||
if cookie.Name == DEVICE_ID_COOKIE && cookie.Value != "" && cookie.Value != y.Addition.DeviceID {
|
||
y.Addition.DeviceID = cookie.Value
|
||
op.MustSaveDriverStorage(y)
|
||
return
|
||
}
|
||
}
|
||
}
|
||
|
||
func (y *Cloud189PC) loginByPassword() (err error) {
|
||
// 初始化登陆所需参数
|
||
if y.loginParam == nil {
|
||
if err = y.initLoginParam(); err != nil {
|
||
// 验证码也通过错误返回
|
||
return err
|
||
}
|
||
}
|
||
// 设备二次校验必须复用同一套登陆参数,此时不能销毁也不能重新初始化
|
||
keepLoginParam := false
|
||
defer func() {
|
||
// 销毁验证码
|
||
y.VCode = ""
|
||
if keepLoginParam {
|
||
y.Status = err.Error()
|
||
op.MustSaveDriverStorage(y)
|
||
return
|
||
}
|
||
// 销毁登陆参数
|
||
y.loginParam = nil
|
||
// 遇到错误,重新加载登陆参数(刷新验证码)
|
||
if err != nil {
|
||
if y.NoUseOcr {
|
||
if err1 := y.initLoginParam(); err1 != nil {
|
||
err = fmt.Errorf("err1: %s \nerr2: %s", err, err1)
|
||
}
|
||
}
|
||
|
||
y.Status = err.Error()
|
||
op.MustSaveDriverStorage(y)
|
||
}
|
||
}()
|
||
|
||
param := y.loginParam
|
||
var loginresp LoginResp
|
||
res, err := y.client.R().
|
||
ForceContentType("application/json;charset=UTF-8").SetResult(&loginresp).
|
||
SetHeaders(y.loginHeaders(param.BaseLoginParam)).
|
||
SetFormData(map[string]string{
|
||
"version": "v2.0",
|
||
"apToken": "",
|
||
"appKey": APP_ID,
|
||
"pageKey": "normal",
|
||
"accountType": ACCOUNT_TYPE,
|
||
"userName": param.RsaUsername,
|
||
"password": param.RsaPassword,
|
||
"epd": param.RsaPassword,
|
||
"validateCode": y.VCode,
|
||
"captchaToken": param.CaptchaToken,
|
||
"returnUrl": RETURN_URL,
|
||
// "mailSuffix": "@189.cn",
|
||
"dynamicCheck": "FALSE",
|
||
"clientType": CLIENT_TYPE,
|
||
"cb_SaveName": "1",
|
||
"isOauth2": "false",
|
||
"state": "",
|
||
"paramId": param.ParamId,
|
||
}).
|
||
Post(AUTH_URL + "/api/logbox/oauth2/loginSubmit.do")
|
||
if err != nil {
|
||
return err
|
||
}
|
||
y.saveDeviceID(res)
|
||
|
||
// 设备二次校验:服务端要求短信验证,保留登陆参数并引导填写短信验证码
|
||
if loginresp.Result == SecondDeviceAuthResult {
|
||
err = y.secondDeviceAuth(loginresp.Mobile)
|
||
// 校验未完成时保留登陆参数,等待用户回填短信验证码
|
||
keepLoginParam = err != nil && y.loginParam != nil
|
||
return err
|
||
}
|
||
|
||
if loginresp.ToUrl == "" {
|
||
return fmt.Errorf("login failed,No toUrl obtained, msg: %s", loginresp.Msg)
|
||
}
|
||
|
||
return y.getSessionByRedirectURL(loginresp.ToUrl)
|
||
}
|
||
|
||
// 设备二次校验:先发短信,用户回填验证码后再提交
|
||
func (y *Cloud189PC) secondDeviceAuth(mobile string) error {
|
||
param := y.loginParam
|
||
if mobile != "" {
|
||
param.SecondAuthMobile = mobile
|
||
}
|
||
if param.SecondAuthMobile == "" {
|
||
return errors.New("second device verification is required, but no mobile was returned")
|
||
}
|
||
|
||
// 已填写短信验证码,直接提交校验
|
||
if y.SmsCode != "" {
|
||
smsCode := y.SmsCode
|
||
y.SmsCode = ""
|
||
op.MustSaveDriverStorage(y)
|
||
return y.submitSecondDeviceAuth(smsCode)
|
||
}
|
||
|
||
var smsResp LoginResp
|
||
_, err := y.client.R().
|
||
ForceContentType("application/json;charset=UTF-8").SetResult(&smsResp).
|
||
SetHeaders(y.loginHeaders(param.BaseLoginParam)).
|
||
SetFormData(map[string]string{
|
||
"mobile": param.SecondAuthMobile,
|
||
"appKey": APP_ID,
|
||
}).
|
||
Post(AUTH_URL + "/api/logbox/oauth2/sendSmsCodeForSecondAuth.do")
|
||
if err != nil {
|
||
return err
|
||
}
|
||
if smsResp.Result != 0 {
|
||
return fmt.Errorf("failed to send the verification SMS: %s", smsResp.Msg)
|
||
}
|
||
// 保留登陆参数,等待用户回填短信验证码后重新保存
|
||
return errors.New("second device verification is required, an SMS code has been sent, please fill it into `sms_code` and save again")
|
||
}
|
||
|
||
// 提交短信验证码完成设备二次校验
|
||
// 注意:该接口没有独立的短信码字段,短信码要加密后放在epd里(登陆时epd装的是密码)
|
||
func (y *Cloud189PC) submitSecondDeviceAuth(smsCode string) error {
|
||
param := y.loginParam
|
||
var authResp LoginResp
|
||
res, err := y.client.R().
|
||
ForceContentType("application/json;charset=UTF-8").SetResult(&authResp).
|
||
SetHeaders(y.loginHeaders(param.BaseLoginParam)).
|
||
SetFormData(map[string]string{
|
||
"mobile": param.SecondAuthMobile,
|
||
"appKey": APP_ID,
|
||
"userName": param.RsaUsername,
|
||
"epd": param.encryptSecret(smsCode),
|
||
"accountType": ACCOUNT_TYPE,
|
||
"returnUrl": RETURN_URL,
|
||
"isOauth2": "false",
|
||
"cb_SaveName": "1",
|
||
"state": "",
|
||
"paramId": param.ParamId,
|
||
}).
|
||
Post(AUTH_URL + "/api/logbox/oauth2/submitForSecondAuth.do")
|
||
if err != nil {
|
||
return err
|
||
}
|
||
// 校验通过后服务端会下发DEVICEID,保存下来以后就不会再触发二次校验
|
||
y.saveDeviceID(res)
|
||
|
||
if authResp.Result != 0 {
|
||
return fmt.Errorf("second device verification failed: %s", authResp.Msg)
|
||
}
|
||
if authResp.ToUrl == "" {
|
||
return fmt.Errorf("second device verification failed, no toUrl obtained, msg: %s", authResp.Msg)
|
||
}
|
||
return y.getSessionByRedirectURL(authResp.ToUrl)
|
||
}
|
||
|
||
// 用登陆结果的跳转地址换取会话
|
||
func (y *Cloud189PC) getSessionByRedirectURL(redirectURL string) error {
|
||
var erron RespErr
|
||
var tokenInfo AppSessionResp
|
||
_, err := y.client.R().
|
||
SetResult(&tokenInfo).SetError(&erron).
|
||
SetQueryParams(clientSuffix()).
|
||
SetQueryParams(y.deviceParams()).
|
||
SetQueryParam("redirectURL", redirectURL).
|
||
SetHeader("X-Request-ID", uuid.NewString()).
|
||
Post(API_URL + "/getSessionForPC.action")
|
||
if err != nil {
|
||
return err
|
||
}
|
||
|
||
if erron.HasError() {
|
||
return &erron
|
||
}
|
||
if tokenInfo.ResCode != 0 {
|
||
return errors.New(tokenInfo.ResMessage)
|
||
}
|
||
y.Addition.AccessToken = tokenInfo.AccessToken
|
||
y.Addition.RefreshToken = tokenInfo.RefreshToken
|
||
y.tokenInfo = &tokenInfo
|
||
op.MustSaveDriverStorage(y)
|
||
return nil
|
||
}
|
||
|
||
func (y *Cloud189PC) loginByQRCode() error {
|
||
if y.qrcodeParam == nil {
|
||
if err := y.initQRCodeParam(); err != nil {
|
||
// 二维码也通过错误返回
|
||
return err
|
||
}
|
||
}
|
||
|
||
// 本地轮询,扫码确认后自动继续,不需要用户反复保存
|
||
deadline := time.Now().Add(QRCODE_POLL_TIMEOUT)
|
||
lastStatus := -106
|
||
for {
|
||
state, err := y.checkQRCodeState()
|
||
if err != nil {
|
||
return fmt.Errorf("failed to check QR code state: %w", err)
|
||
}
|
||
lastStatus = state.Status
|
||
|
||
switch state.Status {
|
||
case 0: // 登录成功
|
||
y.qrcodeParam = nil
|
||
return y.getSessionByRedirectURL(state.RedirectUrl)
|
||
case -106, -11002: // -106 等待扫描,-11002 已扫描等待确认
|
||
case -11001: // 二维码过期
|
||
y.qrcodeParam = nil
|
||
return errors.New("QR code expired, please try again")
|
||
default: // 其他错误
|
||
y.qrcodeParam = nil
|
||
return fmt.Errorf("QR code login failed with status %d: %s", state.Status, state.Msg)
|
||
}
|
||
|
||
if time.Now().Add(QRCODE_POLL_INTERVAL).After(deadline) {
|
||
break
|
||
}
|
||
time.Sleep(QRCODE_POLL_INTERVAL)
|
||
}
|
||
|
||
// 轮询超时,把二维码交回前端等待下一次保存
|
||
if lastStatus == -11002 {
|
||
return y.genQRCode("QR code has been scanned, please confirm the login on your phone and save again")
|
||
}
|
||
return y.genQRCode("QR code has not been scanned yet, please scan and save again")
|
||
}
|
||
|
||
type qrCodeState struct {
|
||
Status int `json:"status"`
|
||
RedirectUrl string `json:"redirectUrl"`
|
||
Msg string `json:"msg"`
|
||
}
|
||
|
||
// 查询扫码状态,参数需与官方PC端一致,否则服务端不会返回授权结果
|
||
func (y *Cloud189PC) checkQRCodeState() (*qrCodeState, error) {
|
||
now := time.Now()
|
||
var state qrCodeState
|
||
_, err := y.client.R().
|
||
SetHeaders(y.loginHeaders(y.qrcodeParam.BaseLoginParam)).
|
||
SetFormData(map[string]string{
|
||
"appId": APP_ID,
|
||
"clientType": QR_CLIENT_TYPE,
|
||
"returnUrl": RETURN_URL,
|
||
"paramId": y.qrcodeParam.ParamId,
|
||
"uuid": y.qrcodeParam.UUID,
|
||
"encryuuid": y.qrcodeParam.EncryUUID,
|
||
"cb_SaveName": "3",
|
||
"isOauth2": "false",
|
||
"state": "",
|
||
"date": formatDate(now),
|
||
"timeStamp": fmt.Sprint(now.UTC().UnixNano() / 1e6),
|
||
}).
|
||
ForceContentType("application/json;charset=UTF-8").
|
||
SetResult(&state).
|
||
Post(AUTH_URL + "/api/logbox/oauth2/qrcodeLoginState.do")
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
return &state, nil
|
||
}
|
||
|
||
func (y *Cloud189PC) genQRCode(text string) error {
|
||
// 展示二维码
|
||
qrTemplate := `<body>
|
||
state: %s
|
||
<br><img src="data:image/jpeg;base64,%s"/>
|
||
<br>Or Click here: <a href="%s">Login</a>
|
||
</body>`
|
||
|
||
// Generate QR code
|
||
qrCode, err := qrcode.Encode(y.qrcodeParam.UUID, qrcode.Medium, 256)
|
||
if err != nil {
|
||
return fmt.Errorf("failed to generate QR code: %v", err)
|
||
}
|
||
|
||
// Encode QR code to base64
|
||
qrCodeBase64 := base64.StdEncoding.EncodeToString(qrCode)
|
||
|
||
// Create the HTML page
|
||
qrPage := fmt.Sprintf(qrTemplate, text, qrCodeBase64, y.qrcodeParam.UUID)
|
||
return fmt.Errorf("need verify: \n%s", qrPage)
|
||
}
|
||
|
||
func (y *Cloud189PC) initBaseParams() (*BaseLoginParam, error) {
|
||
// 清除cookie,并带上已保存的设备标识
|
||
jar, _ := cookiejar.New(nil)
|
||
y.applyDeviceID(jar)
|
||
y.client.SetCookieJar(jar)
|
||
|
||
res, err := y.client.R().
|
||
SetQueryParams(map[string]string{
|
||
"appId": APP_ID,
|
||
"clientType": CLIENT_TYPE,
|
||
"returnURL": RETURN_URL,
|
||
"timeStamp": fmt.Sprint(timestamp()),
|
||
}).
|
||
Get(WEB_URL + "/api/portal/unifyLoginForPC.action")
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
|
||
// 当前登陆页把lt/reqId放在跳转地址上,老页面则写在页内变量里,两种都要支持
|
||
param, err := parseBaseParamFromRedirect(res.RawResponse.Request.URL)
|
||
if err != nil {
|
||
param, err = parseBaseParamFromPage(res.String())
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
return param, nil
|
||
}
|
||
|
||
// 跳转地址上没有paramId,需要再问一次appConf.do
|
||
var appConf AppConfResp
|
||
_, err = y.client.R().
|
||
SetHeaders(y.loginHeaders(*param)).
|
||
ForceContentType("application/json;charset=UTF-8").
|
||
SetResult(&appConf).
|
||
SetFormData(map[string]string{
|
||
"version": "2.0",
|
||
"appKey": APP_ID,
|
||
}).
|
||
Post(AUTH_URL + "/api/logbox/oauth2/appConf.do")
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
if !appConf.Succeeded() || appConf.Data.ParamId == "" {
|
||
return nil, fmt.Errorf("failed to get the login paramId: %s", appConf.Msg)
|
||
}
|
||
param.ParamId = appConf.Data.ParamId
|
||
return param, nil
|
||
}
|
||
|
||
// parseBaseParamFromRedirect 从logbox跳转地址提取登陆参数,并以该地址作为后续请求的Referer
|
||
func parseBaseParamFromRedirect(finalUrl *url.URL) (*BaseLoginParam, error) {
|
||
if finalUrl == nil {
|
||
return nil, errors.New("no login page redirect")
|
||
}
|
||
query := finalUrl.Query()
|
||
lt, reqId := query.Get("lt"), query.Get("reqId")
|
||
if lt == "" || reqId == "" {
|
||
return nil, errors.New("no lt/reqId in the login page redirect")
|
||
}
|
||
return &BaseLoginParam{
|
||
Lt: lt,
|
||
ReqId: reqId,
|
||
Referer: finalUrl.String(),
|
||
}, nil
|
||
}
|
||
|
||
// parseBaseParamFromPage 兼容把参数写在页内变量里的老登陆页
|
||
func parseBaseParamFromPage(body string) (*BaseLoginParam, error) {
|
||
lt, err := matchLoginParam(body, `lt = "(.+?)"`, "lt")
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
reqId, err := matchLoginParam(body, `reqId = "(.+?)"`, "reqId")
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
paramId, err := matchLoginParam(body, `paramId = "(.+?)"`, "paramId")
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
// 老页面才有内嵌的图形验证码token
|
||
captchaToken, _ := matchLoginParam(body, `'captchaToken' value='(.+?)'`, "captchaToken")
|
||
|
||
param := &BaseLoginParam{
|
||
CaptchaToken: captchaToken,
|
||
Lt: lt,
|
||
ParamId: paramId,
|
||
ReqId: reqId,
|
||
}
|
||
encryptUrl, _ := matchLoginParam(body, `encryptUrl = "(.+?)"`, "encryptUrl")
|
||
param.Referer = AUTH_URL + "/api/logbox/separate/web/index.html?" + strings.Join([]string{
|
||
"appId=" + url.QueryEscape(APP_ID),
|
||
"lt=" + url.QueryEscape(param.Lt),
|
||
"reqId=" + url.QueryEscape(param.ReqId),
|
||
}, "&")
|
||
if encryptUrl != "" {
|
||
param.Referer += "&encryptUrl=" + url.QueryEscape(encryptUrl)
|
||
}
|
||
return param, nil
|
||
}
|
||
|
||
// matchLoginParam 从登陆页面提取参数,缺失时返回可读的错误而不是panic
|
||
func matchLoginParam(body, pattern, name string) (string, error) {
|
||
matches := regexp.MustCompile(pattern).FindStringSubmatch(body)
|
||
if len(matches) < 2 {
|
||
return "", fmt.Errorf("failed to get %s from the login page", name)
|
||
}
|
||
return matches[1], nil
|
||
}
|
||
|
||
/* 初始化登陆需要的参数
|
||
* 如果遇到验证码返回错误
|
||
*/
|
||
func (y *Cloud189PC) initLoginParam() error {
|
||
y.loginParam = nil
|
||
|
||
baseParam, err := y.initBaseParams()
|
||
if err != nil {
|
||
return err
|
||
}
|
||
|
||
y.loginParam = &LoginParam{BaseLoginParam: *baseParam}
|
||
|
||
// 获取rsa公钥
|
||
var encryptConf EncryptConfResp
|
||
_, err = y.client.R().
|
||
ForceContentType("application/json;charset=UTF-8").SetResult(&encryptConf).
|
||
SetFormData(map[string]string{"appId": APP_ID}).
|
||
Post(AUTH_URL + "/api/logbox/config/encryptConf.do")
|
||
if err != nil {
|
||
return err
|
||
}
|
||
|
||
y.loginParam.jRsaKey = fmt.Sprintf("-----BEGIN PUBLIC KEY-----\n%s\n-----END PUBLIC KEY-----", encryptConf.Data.PubKey)
|
||
y.loginParam.rsaPrefix = encryptConf.Data.Pre
|
||
y.loginParam.RsaUsername = y.loginParam.encryptSecret(y.Username)
|
||
y.loginParam.RsaPassword = y.loginParam.encryptSecret(y.Password)
|
||
|
||
// 判断是否需要验证码
|
||
resp, err := y.client.R().
|
||
SetHeaders(y.loginHeaders(y.loginParam.BaseLoginParam)).
|
||
SetFormData(map[string]string{
|
||
"appKey": APP_ID,
|
||
"accountType": ACCOUNT_TYPE,
|
||
"userName": y.loginParam.RsaUsername,
|
||
}).Post(AUTH_URL + "/api/logbox/oauth2/needcaptcha.do")
|
||
if err != nil {
|
||
return err
|
||
}
|
||
if resp.String() == "0" {
|
||
return nil
|
||
}
|
||
|
||
// 拉取验证码
|
||
imgRes, err := y.client.R().
|
||
SetQueryParams(map[string]string{
|
||
"token": y.loginParam.CaptchaToken,
|
||
"REQID": y.loginParam.ReqId,
|
||
"rnd": fmt.Sprint(timestamp()),
|
||
}).
|
||
Get(AUTH_URL + "/api/logbox/oauth2/picCaptcha.do")
|
||
if err != nil {
|
||
return fmt.Errorf("failed to obtain verification code")
|
||
}
|
||
if imgRes.Size() > 20 {
|
||
if setting.GetStr(conf.OcrApi) != "" && !y.NoUseOcr {
|
||
vRes, err := base.RestyClient.R().
|
||
SetMultipartField("image", "validateCode.png", "image/png", bytes.NewReader(imgRes.Body())).
|
||
Post(setting.GetStr(conf.OcrApi))
|
||
if err != nil {
|
||
return err
|
||
}
|
||
if jsoniter.Get(vRes.Body(), "status").ToInt() == 200 {
|
||
y.VCode = jsoniter.Get(vRes.Body(), "result").ToString()
|
||
return nil
|
||
}
|
||
}
|
||
|
||
// 返回验证码图片给前端
|
||
return fmt.Errorf(`need img validate code: <img src="data:image/png;base64,%s"/>`, base64.StdEncoding.EncodeToString(imgRes.Body()))
|
||
}
|
||
return nil
|
||
}
|
||
|
||
// getQRCode 获取并返回二维码
|
||
func (y *Cloud189PC) initQRCodeParam() (err error) {
|
||
y.qrcodeParam = nil
|
||
|
||
baseParam, err := y.initBaseParams()
|
||
if err != nil {
|
||
return err
|
||
}
|
||
|
||
var qrcodeParam QRLoginParam
|
||
_, err = y.client.R().
|
||
SetHeaders(y.loginHeaders(*baseParam)).
|
||
SetFormData(map[string]string{"appId": APP_ID}).
|
||
ForceContentType("application/json;charset=UTF-8").
|
||
SetResult(&qrcodeParam).
|
||
Post(AUTH_URL + "/api/logbox/oauth2/getUUID.do")
|
||
if err != nil {
|
||
return err
|
||
}
|
||
if qrcodeParam.UUID == "" {
|
||
return errors.New("failed to get the QR code uuid")
|
||
}
|
||
qrcodeParam.BaseLoginParam = *baseParam
|
||
y.qrcodeParam = &qrcodeParam
|
||
|
||
return y.genQRCode("please scan the QR code with the 189 Cloud app, then save the settings again.")
|
||
}
|
||
|
||
// 刷新会话
|
||
func (y *Cloud189PC) refreshSession() (err error) {
|
||
return y.refreshSessionWithRetry(0)
|
||
}
|
||
|
||
func (y *Cloud189PC) refreshSessionWithRetry(retryCount int) (err error) {
|
||
if y.ref != nil {
|
||
return y.ref.refreshSessionWithRetry(retryCount)
|
||
}
|
||
var erron RespErr
|
||
var userSessionResp UserSessionResp
|
||
_, err = y.client.R().
|
||
SetResult(&userSessionResp).SetError(&erron).
|
||
SetQueryParams(clientSuffix()).
|
||
SetQueryParams(y.deviceParams()).
|
||
SetQueryParams(map[string]string{
|
||
"appId": APP_ID,
|
||
"accessToken": y.tokenInfo.AccessToken,
|
||
}).
|
||
SetHeader("X-Request-ID", uuid.NewString()).
|
||
Get(API_URL + "/getSessionForPC.action")
|
||
if err != nil {
|
||
return err
|
||
}
|
||
|
||
// token生效刷新token
|
||
if erron.HasError() {
|
||
if erron.ResCode == UserInvalidOpenTokenError {
|
||
return y.refreshTokenWithRetry(retryCount)
|
||
}
|
||
return &erron
|
||
}
|
||
y.tokenInfo.UserSessionResp = userSessionResp
|
||
return nil
|
||
}
|
||
|
||
// refreshToken 刷新token,失败时返回错误,不再直接调用login
|
||
func (y *Cloud189PC) refreshToken() (err error) {
|
||
return y.refreshTokenWithRetry(0)
|
||
}
|
||
|
||
func (y *Cloud189PC) refreshTokenWithRetry(retryCount int) (err error) {
|
||
if y.ref != nil {
|
||
return y.ref.refreshTokenWithRetry(retryCount)
|
||
}
|
||
|
||
// 限制重试次数,避免无限递归
|
||
if retryCount >= 3 {
|
||
if y.Addition.RefreshToken != "" {
|
||
y.Addition.RefreshToken = ""
|
||
op.MustSaveDriverStorage(y)
|
||
}
|
||
return errors.New("refresh token failed after maximum retries")
|
||
}
|
||
|
||
// 该接口刷新失败时以HTTP 200返回 result/msg,SetError不会触发,必须解析响应体判断
|
||
var tokenInfo RefreshTokenResp
|
||
_, err = y.client.R().
|
||
SetResult(&tokenInfo).
|
||
ForceContentType("application/json;charset=UTF-8").
|
||
SetFormData(map[string]string{
|
||
"clientId": APP_ID,
|
||
"refreshToken": y.tokenInfo.RefreshToken,
|
||
"grantType": "refresh_token",
|
||
"format": "json",
|
||
}).
|
||
Post(AUTH_URL + "/api/oauth2/refreshToken.do")
|
||
if err != nil {
|
||
return err
|
||
}
|
||
|
||
// 如果刷新失败,返回错误给上层处理
|
||
if tokenInfo.HasError() {
|
||
refreshErr := tokenInfo.Error()
|
||
if y.Addition.RefreshToken != "" {
|
||
y.Addition.RefreshToken = ""
|
||
op.MustSaveDriverStorage(y)
|
||
}
|
||
|
||
// 根据登录类型决定下一步行为
|
||
if y.LoginType == "qrcode" {
|
||
return fmt.Errorf("QR code session has expired, please re-scan the code to log in: %s", refreshErr)
|
||
}
|
||
// 没有账号密码时无法回退到完整登录,直接把刷新失败的原因返回
|
||
if y.Username == "" || y.Password == "" {
|
||
return errors.New(refreshErr)
|
||
}
|
||
// 密码登录模式下,尝试回退到完整登录
|
||
return y.login()
|
||
}
|
||
|
||
y.Addition.AccessToken = tokenInfo.AccessToken
|
||
y.Addition.RefreshToken = tokenInfo.RefreshToken
|
||
y.tokenInfo.AccessToken = tokenInfo.AccessToken
|
||
y.tokenInfo.RefreshToken = tokenInfo.RefreshToken
|
||
op.MustSaveDriverStorage(y)
|
||
return y.refreshSessionWithRetry(retryCount + 1)
|
||
}
|
||
|
||
func (y *Cloud189PC) keepAlive() {
|
||
_, err := y.get(API_URL+"/keepUserSession.action", func(r *resty.Request) {
|
||
r.SetQueryParams(clientSuffix())
|
||
}, nil)
|
||
if err != nil {
|
||
utils.Log.Warnf("189pc: Failed to keep user session alive: %v", err)
|
||
// 如果keepAlive失败,尝试刷新session
|
||
if refreshErr := y.refreshSession(); refreshErr != nil {
|
||
utils.Log.Errorf("189pc: Failed to refresh session after keepAlive error: %v", refreshErr)
|
||
}
|
||
} else {
|
||
utils.Log.Debugf("189pc: User session kept alive successfully.")
|
||
}
|
||
}
|
||
|
||
// 普通上传
|
||
// 无法上传大小为0的文件
|
||
func (y *Cloud189PC) StreamUpload(ctx context.Context, dstDir model.Obj, file model.FileStreamer, up driver.UpdateProgress, isFamily bool, overwrite bool) (model.Obj, error) {
|
||
// 文件大小
|
||
fileSize := file.GetSize()
|
||
// 分片大小,不得为文件大小
|
||
sliceSize := partSize(fileSize)
|
||
|
||
params := Params{
|
||
"parentFolderId": dstDir.GetID(),
|
||
"fileName": url.QueryEscape(file.GetName()),
|
||
"fileSize": fmt.Sprint(fileSize),
|
||
"sliceSize": fmt.Sprint(sliceSize), // 必须为特定分片大小
|
||
"lazyCheck": "1",
|
||
}
|
||
|
||
fullUrl := UPLOAD_URL
|
||
if isFamily {
|
||
params.Set("familyId", y.FamilyID)
|
||
fullUrl += "/family"
|
||
} else {
|
||
// params.Set("extend", `{"opScene":"1","relativepath":"","rootfolderid":""}`)
|
||
fullUrl += "/person"
|
||
}
|
||
|
||
// 初始化上传
|
||
var initMultiUpload InitMultiUploadResp
|
||
_, err := y.request(fullUrl+"/initMultiUpload", http.MethodGet, func(req *resty.Request) {
|
||
req.SetContext(ctx)
|
||
}, params, &initMultiUpload, isFamily)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
|
||
ss, err := stream.NewStreamSectionReader(file, int(sliceSize), &up)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
|
||
threadG, upCtx := errgroup.NewOrderedGroupWithContext(ctx, y.uploadThread,
|
||
retry.Attempts(3),
|
||
retry.Delay(time.Second),
|
||
retry.DelayType(retry.BackOffDelay))
|
||
|
||
count := 1
|
||
if fileSize > sliceSize {
|
||
count = int((fileSize + sliceSize - 1) / sliceSize)
|
||
}
|
||
lastPartSize := fileSize % sliceSize
|
||
if lastPartSize == 0 {
|
||
lastPartSize = sliceSize
|
||
}
|
||
|
||
silceMd5Hexs := make([]string, 0, count)
|
||
silceMd5 := utils.MD5.NewFunc()
|
||
var writers io.Writer = silceMd5
|
||
|
||
// 如果启用了 torrent 生成,额外计算 SHA-1 piece hash
|
||
generateTorrent := y.Addition.GenerateTorrent
|
||
pieceSHA1Hashes := make([]byte, 0, count*20)
|
||
|
||
fileMd5Hex := file.GetHash().GetHash(utils.MD5)
|
||
var fileMd5 hash.Hash
|
||
if len(fileMd5Hex) != utils.MD5.Width {
|
||
fileMd5 = utils.MD5.NewFunc()
|
||
writers = io.MultiWriter(silceMd5, fileMd5)
|
||
}
|
||
for i := 1; i <= count; i++ {
|
||
if utils.IsCanceled(upCtx) {
|
||
break
|
||
}
|
||
offset := int64((i)-1) * sliceSize
|
||
partSize := sliceSize
|
||
if i == count {
|
||
partSize = lastPartSize
|
||
}
|
||
partInfo := ""
|
||
var reader io.ReadSeeker
|
||
threadG.GoWithLifecycle(errgroup.Lifecycle{
|
||
Before: func(ctx context.Context) (err error) {
|
||
reader, err = ss.GetSectionReader(offset, partSize)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
silceMd5.Reset()
|
||
|
||
// 如果需要生成 torrent,同时计算 SHA-1
|
||
var sha1Writer hash.Hash
|
||
var multiWriter io.Writer
|
||
if generateTorrent {
|
||
sha1Writer = sha1Pkg.New()
|
||
multiWriter = io.MultiWriter(writers, sha1Writer)
|
||
} else {
|
||
multiWriter = writers
|
||
}
|
||
|
||
w, err := utils.CopyWithBuffer(multiWriter, reader)
|
||
if w != partSize {
|
||
return fmt.Errorf("failed to read all data: (expect =%d, actual =%d) %w", partSize, w, err)
|
||
}
|
||
// 计算块md5并进行hex和base64编码
|
||
md5Bytes := silceMd5.Sum(nil)
|
||
silceMd5Hexs = append(silceMd5Hexs, strings.ToUpper(hex.EncodeToString(md5Bytes)))
|
||
partInfo = fmt.Sprintf("%d-%s", i, base64.StdEncoding.EncodeToString(md5Bytes))
|
||
|
||
// 收集 SHA-1 piece hash
|
||
if generateTorrent && sha1Writer != nil {
|
||
pieceSHA1Hashes = append(pieceSHA1Hashes, sha1Writer.Sum(nil)...)
|
||
}
|
||
return nil
|
||
},
|
||
Do: func(ctx context.Context) (err error) {
|
||
reader.Seek(0, io.SeekStart)
|
||
uploadUrls, err := y.GetMultiUploadUrls(ctx, isFamily, initMultiUpload.Data.UploadFileID, partInfo)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
|
||
// step.4 上传切片
|
||
uploadUrl := uploadUrls[0]
|
||
_, err = y.put(ctx, uploadUrl.RequestURL, uploadUrl.Headers, false, driver.NewLimitedUploadStream(ctx, reader), isFamily)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
up(float64(threadG.Success()+1) * 100 / float64(count+1))
|
||
return nil
|
||
},
|
||
After: func(err error) {
|
||
ss.FreeSectionReader(reader)
|
||
},
|
||
},
|
||
)
|
||
}
|
||
if err = threadG.Wait(); err != nil {
|
||
return nil, err
|
||
}
|
||
defer up(100)
|
||
|
||
if fileMd5 != nil {
|
||
fileMd5Hex = strings.ToUpper(hex.EncodeToString(fileMd5.Sum(nil)))
|
||
}
|
||
sliceMd5Hex := fileMd5Hex
|
||
if fileSize > sliceSize {
|
||
sliceMd5Hex = strings.ToUpper(utils.GetMD5EncodeStr(strings.Join(silceMd5Hexs, "\n")))
|
||
}
|
||
|
||
// 提交上传
|
||
var resp CommitMultiUploadFileResp
|
||
_, err = y.request(fullUrl+"/commitMultiUploadFile", http.MethodGet,
|
||
func(req *resty.Request) {
|
||
req.SetContext(ctx)
|
||
}, Params{
|
||
"uploadFileId": initMultiUpload.Data.UploadFileID,
|
||
"fileMd5": fileMd5Hex,
|
||
"sliceMd5": sliceMd5Hex,
|
||
"lazyCheck": "1",
|
||
"isLog": "0",
|
||
"opertype": IF(overwrite, "3", "1"),
|
||
}, &resp, isFamily)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
|
||
// 生成 torrent 文件(异步,不影响上传结果)
|
||
if generateTorrent && len(pieceSHA1Hashes) > 0 {
|
||
// 捕获必要的变量
|
||
capturedDstDir := dstDir
|
||
capturedIsFamily := isFamily
|
||
capturedFileName := file.GetName()
|
||
go func() {
|
||
torrentData, err := GenerateTorrent(capturedFileName, fileSize, fileMd5Hex, silceMd5Hexs, sliceSize, pieceSHA1Hashes)
|
||
if err != nil {
|
||
utils.Log.Warnf("生成 torrent 失败: %v", err)
|
||
return
|
||
}
|
||
infoHash, _ := GetInfoHashHex(torrentData)
|
||
torrentName := capturedFileName + ".cas.torrent"
|
||
utils.Log.Infof("已生成 torrent: %s (info_hash: %s, size: %d bytes)",
|
||
torrentName, infoHash, len(torrentData))
|
||
|
||
// 将 torrent 文件上传到同一目录(使用 FastUpload,因为 torrent 文件很小)
|
||
torrentFileStream := &stream.FileStream{
|
||
Ctx: context.Background(),
|
||
Obj: &model.Object{
|
||
Name: torrentName,
|
||
Size: int64(len(torrentData)),
|
||
IsFolder: false,
|
||
},
|
||
Reader: bytes.NewReader(torrentData),
|
||
Mimetype: "application/x-bittorrent",
|
||
}
|
||
_, uploadErr := y.FastUpload(context.Background(), capturedDstDir, torrentFileStream, func(p float64) {}, capturedIsFamily, false)
|
||
if uploadErr != nil {
|
||
utils.Log.Warnf("上传 torrent 文件失败: %v", uploadErr)
|
||
} else {
|
||
utils.Log.Infof("torrent 文件已上传: %s", torrentName)
|
||
op.Cache.DeleteDirectory(y, capturedDstDir.GetPath())
|
||
}
|
||
}()
|
||
}
|
||
|
||
return resp.toFile(), nil
|
||
}
|
||
|
||
func (y *Cloud189PC) RapidUpload(ctx context.Context, dstDir model.Obj, stream model.FileStreamer, isFamily bool, overwrite bool) (model.Obj, error) {
|
||
fileMd5 := stream.GetHash().GetHash(utils.MD5)
|
||
if len(fileMd5) < utils.MD5.Width {
|
||
return nil, errors.New("invalid hash")
|
||
}
|
||
|
||
uploadInfo, err := y.OldUploadCreate(ctx, dstDir.GetID(), fileMd5, stream.GetName(), fmt.Sprint(stream.GetSize()), isFamily)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
|
||
if uploadInfo.FileDataExists != 1 {
|
||
return nil, errors.New("rapid upload fail")
|
||
}
|
||
|
||
return y.OldUploadCommit(ctx, uploadInfo.FileCommitUrl, uploadInfo.UploadFileId, isFamily, overwrite)
|
||
}
|
||
|
||
// 快传
|
||
func (y *Cloud189PC) FastUpload(ctx context.Context, dstDir model.Obj, file model.FileStreamer, up driver.UpdateProgress, isFamily bool, overwrite bool) (model.Obj, error) {
|
||
generateTorrent := y.Addition.GenerateTorrent && !isCASTorrentFile(file.GetName())
|
||
return y.fastUpload(ctx, dstDir, file, up, isFamily, overwrite, generateTorrent)
|
||
}
|
||
|
||
func (y *Cloud189PC) fastUpload(ctx context.Context, dstDir model.Obj, file model.FileStreamer, up driver.UpdateProgress, isFamily bool, overwrite bool, generateTorrent bool) (model.Obj, error) {
|
||
var (
|
||
cache = file.GetFile()
|
||
tmpF *os.File
|
||
err error
|
||
)
|
||
size := file.GetSize()
|
||
if _, ok := cache.(io.ReaderAt); !ok && size > 0 {
|
||
tmpF, err = os.CreateTemp(conf.Conf.TempDir, "file-*")
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
defer func() {
|
||
_ = tmpF.Close()
|
||
_ = os.Remove(tmpF.Name())
|
||
}()
|
||
cache = tmpF
|
||
}
|
||
sliceSize := partSize(size)
|
||
count := 1
|
||
if size > sliceSize {
|
||
count = int((size + sliceSize - 1) / sliceSize)
|
||
}
|
||
lastSliceSize := size % sliceSize
|
||
if lastSliceSize == 0 {
|
||
lastSliceSize = sliceSize
|
||
}
|
||
|
||
// step.1 优先计算所需信息
|
||
byteSize := sliceSize
|
||
fileMd5 := utils.MD5.NewFunc()
|
||
sliceMd5 := utils.MD5.NewFunc()
|
||
sliceMd5Hexs := make([]string, 0, count)
|
||
partInfos := make([]string, 0, count)
|
||
writers := []io.Writer{fileMd5, sliceMd5}
|
||
if tmpF != nil {
|
||
writers = append(writers, tmpF)
|
||
}
|
||
|
||
pieceSHA1Hashes := make([]byte, 0, count*20)
|
||
|
||
written := int64(0)
|
||
for i := 1; i <= count; i++ {
|
||
if utils.IsCanceled(ctx) {
|
||
return nil, ctx.Err()
|
||
}
|
||
|
||
if i == count {
|
||
byteSize = lastSliceSize
|
||
}
|
||
|
||
// 如果需要生成 torrent,同时计算 SHA-1
|
||
var sha1Writer hash.Hash
|
||
var multiWriter io.Writer
|
||
if generateTorrent {
|
||
sha1Writer = sha1Pkg.New()
|
||
multiWriter = io.MultiWriter(append(writers, sha1Writer)...)
|
||
} else {
|
||
multiWriter = io.MultiWriter(writers...)
|
||
}
|
||
|
||
n, err := utils.CopyWithBufferN(multiWriter, file, byteSize)
|
||
written += n
|
||
if err != nil && err != io.EOF {
|
||
return nil, err
|
||
}
|
||
md5Byte := sliceMd5.Sum(nil)
|
||
sliceMd5Hexs = append(sliceMd5Hexs, strings.ToUpper(hex.EncodeToString(md5Byte)))
|
||
partInfos = append(partInfos, fmt.Sprint(i, "-", base64.StdEncoding.EncodeToString(md5Byte)))
|
||
sliceMd5.Reset()
|
||
|
||
// 收集 SHA-1 piece hash(仅在本次分片实际写入了数据时追加)
|
||
if generateTorrent && n > 0 {
|
||
pieceSHA1Hashes = append(pieceSHA1Hashes, sha1Writer.Sum(nil)...)
|
||
}
|
||
}
|
||
|
||
if tmpF != nil {
|
||
if size > 0 && written != size {
|
||
return nil, errs.NewErr(err, "CreateTempFile failed, incoming stream actual size= %d, expect = %d ", written, size)
|
||
}
|
||
_, err = tmpF.Seek(0, io.SeekStart)
|
||
if err != nil {
|
||
return nil, errs.NewErr(err, "CreateTempFile failed, can't seek to 0 ")
|
||
}
|
||
}
|
||
|
||
fileMd5Hex := strings.ToUpper(hex.EncodeToString(fileMd5.Sum(nil)))
|
||
sliceMd5Hex := fileMd5Hex
|
||
if size > sliceSize {
|
||
sliceMd5Hex = strings.ToUpper(utils.GetMD5EncodeStr(strings.Join(sliceMd5Hexs, "\n")))
|
||
}
|
||
|
||
fullUrl := UPLOAD_URL
|
||
if isFamily {
|
||
fullUrl += "/family"
|
||
} else {
|
||
// params.Set("extend", `{"opScene":"1","relativepath":"","rootfolderid":""}`)
|
||
fullUrl += "/person"
|
||
}
|
||
|
||
// 尝试恢复进度
|
||
uploadProgress, ok := base.GetUploadProgress[*UploadProgress](y, y.getTokenInfo().SessionKey, fileMd5Hex)
|
||
if !ok {
|
||
// step.2 预上传
|
||
params := Params{
|
||
"parentFolderId": dstDir.GetID(),
|
||
"fileName": url.QueryEscape(file.GetName()),
|
||
"fileSize": fmt.Sprint(file.GetSize()),
|
||
"fileMd5": fileMd5Hex,
|
||
"sliceSize": fmt.Sprint(sliceSize),
|
||
"sliceMd5": sliceMd5Hex,
|
||
}
|
||
if isFamily {
|
||
params.Set("familyId", y.FamilyID)
|
||
}
|
||
var uploadInfo InitMultiUploadResp
|
||
_, err = y.request(fullUrl+"/initMultiUpload", http.MethodGet, func(req *resty.Request) {
|
||
req.SetContext(ctx)
|
||
}, params, &uploadInfo, isFamily)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
uploadProgress = &UploadProgress{
|
||
UploadInfo: uploadInfo,
|
||
UploadParts: partInfos,
|
||
}
|
||
}
|
||
|
||
uploadInfo := uploadProgress.UploadInfo.Data
|
||
// 网盘中不存在该文件,开始上传
|
||
if uploadInfo.FileDataExists != 1 {
|
||
threadG, upCtx := errgroup.NewGroupWithContext(ctx, y.uploadThread,
|
||
retry.Attempts(3),
|
||
retry.Delay(time.Second),
|
||
retry.DelayType(retry.BackOffDelay))
|
||
for i, uploadPart := range uploadProgress.UploadParts {
|
||
if utils.IsCanceled(upCtx) {
|
||
break
|
||
}
|
||
|
||
i, uploadPart := i, uploadPart
|
||
threadG.Go(func(ctx context.Context) error {
|
||
// step.3 获取上传链接
|
||
uploadUrls, err := y.GetMultiUploadUrls(ctx, isFamily, uploadInfo.UploadFileID, uploadPart)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
uploadUrl := uploadUrls[0]
|
||
|
||
byteSize, offset := sliceSize, int64(uploadUrl.PartNumber-1)*sliceSize
|
||
if uploadUrl.PartNumber == count {
|
||
byteSize = lastSliceSize
|
||
}
|
||
|
||
// step.4 上传切片
|
||
rateLimitedRd := driver.NewLimitedUploadStream(ctx, io.NewSectionReader(cache, offset, byteSize))
|
||
_, err = y.put(ctx, uploadUrl.RequestURL, uploadUrl.Headers, false, rateLimitedRd, isFamily)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
|
||
up(float64(threadG.Success()+1) * 100 / float64(len(uploadUrls)+1))
|
||
uploadProgress.UploadParts[i] = ""
|
||
return nil
|
||
})
|
||
}
|
||
if err = threadG.Wait(); err != nil {
|
||
if errors.Is(err, context.Canceled) {
|
||
uploadProgress.UploadParts = utils.SliceFilter(uploadProgress.UploadParts, func(s string) bool { return s != "" })
|
||
base.SaveUploadProgress(y, uploadProgress, y.getTokenInfo().SessionKey, fileMd5Hex)
|
||
}
|
||
return nil, err
|
||
}
|
||
defer up(100)
|
||
}
|
||
|
||
// step.5 提交
|
||
var resp CommitMultiUploadFileResp
|
||
_, err = y.request(fullUrl+"/commitMultiUploadFile", http.MethodGet,
|
||
func(req *resty.Request) {
|
||
req.SetContext(ctx)
|
||
}, Params{
|
||
"uploadFileId": uploadInfo.UploadFileID,
|
||
"isLog": "0",
|
||
"opertype": IF(overwrite, "3", "1"),
|
||
}, &resp, isFamily)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
|
||
// 生成 torrent 文件(异步,不影响上传结果)
|
||
if generateTorrent && size > 0 && len(pieceSHA1Hashes) > 0 {
|
||
capturedDstDir := dstDir
|
||
capturedIsFamily := isFamily
|
||
capturedFileName := file.GetName()
|
||
go func() {
|
||
torrentData, err := GenerateTorrent(capturedFileName, size, fileMd5Hex, sliceMd5Hexs, sliceSize, pieceSHA1Hashes)
|
||
if err != nil {
|
||
utils.Log.Warnf("生成 torrent 失败: %v", err)
|
||
return
|
||
}
|
||
infoHash, _ := GetInfoHashHex(torrentData)
|
||
torrentName := capturedFileName + ".cas.torrent"
|
||
utils.Log.Infof("已生成 torrent: %s (info_hash: %s, size: %d bytes)",
|
||
torrentName, infoHash, len(torrentData))
|
||
|
||
// 将 torrent 文件上传到同一目录
|
||
torrentFileStream := &stream.FileStream{
|
||
Ctx: context.Background(),
|
||
Obj: &model.Object{
|
||
Name: torrentName,
|
||
Size: int64(len(torrentData)),
|
||
IsFolder: false,
|
||
},
|
||
Reader: bytes.NewReader(torrentData),
|
||
Mimetype: "application/x-bittorrent",
|
||
}
|
||
_, uploadErr := y.fastUpload(context.Background(), capturedDstDir, torrentFileStream, func(p float64) {}, capturedIsFamily, false, false)
|
||
if uploadErr != nil {
|
||
utils.Log.Warnf("上传 torrent 文件失败: %v", uploadErr)
|
||
} else {
|
||
utils.Log.Infof("torrent 文件已上传: %s", torrentName)
|
||
op.Cache.DeleteDirectory(y, capturedDstDir.GetPath())
|
||
}
|
||
}()
|
||
}
|
||
|
||
return resp.toFile(), nil
|
||
}
|
||
|
||
// 获取上传切片信息
|
||
// 对http body有大小限制,分片信息太多会出错
|
||
func (y *Cloud189PC) GetMultiUploadUrls(ctx context.Context, isFamily bool, uploadFileId string, partInfo ...string) ([]UploadUrlInfo, error) {
|
||
fullUrl := UPLOAD_URL
|
||
if isFamily {
|
||
fullUrl += "/family"
|
||
} else {
|
||
fullUrl += "/person"
|
||
}
|
||
|
||
var uploadUrlsResp UploadUrlsResp
|
||
_, err := y.request(fullUrl+"/getMultiUploadUrls", http.MethodGet,
|
||
func(req *resty.Request) {
|
||
req.SetContext(ctx)
|
||
}, Params{
|
||
"uploadFileId": uploadFileId,
|
||
"partInfo": strings.Join(partInfo, ","),
|
||
}, &uploadUrlsResp, isFamily)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
uploadUrls := uploadUrlsResp.Data
|
||
|
||
if len(uploadUrls) != len(partInfo) {
|
||
return nil, fmt.Errorf("uploadUrls get error, due to get length %d, real length %d", len(partInfo), len(uploadUrls))
|
||
}
|
||
|
||
uploadUrlInfos := make([]UploadUrlInfo, 0, len(uploadUrls))
|
||
for k, uploadUrl := range uploadUrls {
|
||
partNumber, err := strconv.Atoi(strings.TrimPrefix(k, "partNumber_"))
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
uploadUrlInfos = append(uploadUrlInfos, UploadUrlInfo{
|
||
PartNumber: partNumber,
|
||
Headers: ParseHttpHeader(uploadUrl.RequestHeader),
|
||
UploadUrlsData: uploadUrl,
|
||
})
|
||
}
|
||
sort.Slice(uploadUrlInfos, func(i, j int) bool {
|
||
return uploadUrlInfos[i].PartNumber < uploadUrlInfos[j].PartNumber
|
||
})
|
||
return uploadUrlInfos, nil
|
||
}
|
||
|
||
// 旧版本上传,家庭云不支持覆盖
|
||
func (y *Cloud189PC) OldUpload(ctx context.Context, dstDir model.Obj, file model.FileStreamer, up driver.UpdateProgress, isFamily bool, overwrite bool) (model.Obj, error) {
|
||
tempFile, fileMd5, err := stream.CacheFullAndHash(file, &up, utils.MD5)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
rateLimited := driver.NewLimitedUploadStream(ctx, io.NopCloser(tempFile))
|
||
|
||
// 创建上传会话
|
||
uploadInfo, err := y.OldUploadCreate(ctx, dstDir.GetID(), fileMd5, file.GetName(), fmt.Sprint(file.GetSize()), isFamily)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
|
||
// 网盘中不存在该文件,开始上传
|
||
status := GetUploadFileStatusResp{CreateUploadFileResp: *uploadInfo}
|
||
for status.GetSize() < file.GetSize() && status.FileDataExists != 1 {
|
||
if utils.IsCanceled(ctx) {
|
||
return nil, ctx.Err()
|
||
}
|
||
|
||
header := map[string]string{
|
||
"ResumePolicy": "1",
|
||
"Expect": "100-continue",
|
||
}
|
||
|
||
if isFamily {
|
||
header["FamilyId"] = fmt.Sprint(y.FamilyID)
|
||
header["UploadFileId"] = fmt.Sprint(status.UploadFileId)
|
||
} else {
|
||
header["Edrive-UploadFileId"] = fmt.Sprint(status.UploadFileId)
|
||
}
|
||
|
||
_, err := y.put(ctx, status.FileUploadUrl, header, true, rateLimited, isFamily)
|
||
if err, ok := err.(*RespErr); ok && err.Code != "InputStreamReadError" {
|
||
return nil, err
|
||
}
|
||
|
||
// 获取断点状态
|
||
fullUrl := API_URL + "/getUploadFileStatus.action"
|
||
if y.isFamily() {
|
||
fullUrl = API_URL + "/family/file/getFamilyFileStatus.action"
|
||
}
|
||
_, err = y.get(fullUrl, func(req *resty.Request) {
|
||
req.SetContext(ctx).SetQueryParams(map[string]string{
|
||
"uploadFileId": fmt.Sprint(status.UploadFileId),
|
||
"resumePolicy": "1",
|
||
})
|
||
if isFamily {
|
||
req.SetQueryParam("familyId", fmt.Sprint(y.FamilyID))
|
||
}
|
||
}, &status, isFamily)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
if _, err := tempFile.Seek(status.GetSize(), io.SeekStart); err != nil {
|
||
return nil, err
|
||
}
|
||
up(float64(status.GetSize()) / float64(file.GetSize()) * 100)
|
||
}
|
||
|
||
return y.OldUploadCommit(ctx, status.FileCommitUrl, status.UploadFileId, isFamily, overwrite)
|
||
}
|
||
|
||
// 创建上传会话
|
||
func (y *Cloud189PC) OldUploadCreate(ctx context.Context, parentID string, fileMd5, fileName, fileSize string, isFamily bool) (*CreateUploadFileResp, error) {
|
||
var uploadInfo CreateUploadFileResp
|
||
|
||
fullUrl := API_URL + "/createUploadFile.action"
|
||
if isFamily {
|
||
fullUrl = API_URL + "/family/file/createFamilyFile.action"
|
||
}
|
||
_, err := y.post(fullUrl, func(req *resty.Request) {
|
||
req.SetContext(ctx)
|
||
if isFamily {
|
||
req.SetQueryParams(map[string]string{
|
||
"familyId": y.FamilyID,
|
||
"parentId": parentID,
|
||
"fileMd5": fileMd5,
|
||
"fileName": fileName,
|
||
"fileSize": fileSize,
|
||
"resumePolicy": "1",
|
||
})
|
||
} else {
|
||
req.SetFormData(map[string]string{
|
||
"parentFolderId": parentID,
|
||
"fileName": fileName,
|
||
"size": fileSize,
|
||
"md5": fileMd5,
|
||
"opertype": "3",
|
||
"flag": "1",
|
||
"resumePolicy": "1",
|
||
"isLog": "0",
|
||
})
|
||
}
|
||
}, &uploadInfo, isFamily)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
return &uploadInfo, nil
|
||
}
|
||
|
||
// 提交上传文件
|
||
func (y *Cloud189PC) OldUploadCommit(ctx context.Context, fileCommitUrl string, uploadFileID int64, isFamily bool, overwrite bool) (model.Obj, error) {
|
||
var resp OldCommitUploadFileResp
|
||
_, err := y.post(fileCommitUrl, func(req *resty.Request) {
|
||
req.SetContext(ctx)
|
||
if isFamily {
|
||
req.SetHeaders(map[string]string{
|
||
"ResumePolicy": "1",
|
||
"UploadFileId": fmt.Sprint(uploadFileID),
|
||
"FamilyId": fmt.Sprint(y.FamilyID),
|
||
})
|
||
} else {
|
||
req.SetFormData(map[string]string{
|
||
"opertype": IF(overwrite, "3", "1"),
|
||
"resumePolicy": "1",
|
||
"uploadFileId": fmt.Sprint(uploadFileID),
|
||
"isLog": "0",
|
||
})
|
||
}
|
||
}, &resp, isFamily)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
return resp.toFile(), nil
|
||
}
|
||
|
||
func (y *Cloud189PC) isFamily() bool {
|
||
return y.Type == "family"
|
||
}
|
||
|
||
func (y *Cloud189PC) isLogin() bool {
|
||
if y.tokenInfo == nil {
|
||
return false
|
||
}
|
||
_, err := y.get(API_URL+"/getUserInfo.action", nil, nil)
|
||
return err == nil
|
||
}
|
||
|
||
// 创建家庭云中转文件夹
|
||
func (y *Cloud189PC) createFamilyTransferFolder() error {
|
||
var rootFolder Cloud189Folder
|
||
_, err := y.post(API_URL+"/family/file/createFolder.action", func(req *resty.Request) {
|
||
req.SetQueryParams(map[string]string{
|
||
"folderName": "FamilyTransferFolder",
|
||
"familyId": y.FamilyID,
|
||
})
|
||
}, &rootFolder, true)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
y.familyTransferFolder = &rootFolder
|
||
return nil
|
||
}
|
||
|
||
// 清理中转文件夹
|
||
func (y *Cloud189PC) cleanFamilyTransfer(ctx context.Context) error {
|
||
transferFolderId := y.familyTransferFolder.GetID()
|
||
for pageNum := 1; ; pageNum++ {
|
||
resp, err := y.getFilesWithPage(ctx, transferFolderId, true, pageNum, 100, "lastOpTime", "asc")
|
||
if err != nil {
|
||
return err
|
||
}
|
||
// 获取完毕跳出
|
||
if resp.FileListAO.Count == 0 {
|
||
break
|
||
}
|
||
|
||
var tasks []BatchTaskInfo
|
||
for i := 0; i < len(resp.FileListAO.FolderList); i++ {
|
||
folder := resp.FileListAO.FolderList[i]
|
||
tasks = append(tasks, BatchTaskInfo{
|
||
FileId: folder.GetID(),
|
||
FileName: folder.GetName(),
|
||
IsFolder: BoolToNumber(folder.IsDir()),
|
||
})
|
||
}
|
||
for i := 0; i < len(resp.FileListAO.FileList); i++ {
|
||
file := resp.FileListAO.FileList[i]
|
||
tasks = append(tasks, BatchTaskInfo{
|
||
FileId: file.GetID(),
|
||
FileName: file.GetName(),
|
||
IsFolder: BoolToNumber(file.IsDir()),
|
||
})
|
||
}
|
||
|
||
if len(tasks) > 0 {
|
||
// 删除
|
||
resp, err := y.CreateBatchTask("DELETE", y.FamilyID, "", nil, tasks...)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
err = y.WaitBatchTask("DELETE", resp.TaskID, time.Second)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
// 永久删除
|
||
resp, err = y.CreateBatchTask("CLEAR_RECYCLE", y.FamilyID, "", nil, tasks...)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
err = y.WaitBatchTask("CLEAR_RECYCLE", resp.TaskID, time.Second)
|
||
return err
|
||
}
|
||
}
|
||
return nil
|
||
}
|
||
|
||
// 获取家庭云所有用户信息
|
||
func (y *Cloud189PC) getFamilyInfoList() ([]FamilyInfoResp, error) {
|
||
var resp FamilyInfoListResp
|
||
_, err := y.get(API_URL+"/family/manage/getFamilyList.action", nil, &resp, true)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
return resp.FamilyInfoResp, nil
|
||
}
|
||
|
||
// 抽取家庭云ID
|
||
func (y *Cloud189PC) getFamilyID() (string, error) {
|
||
infos, err := y.getFamilyInfoList()
|
||
if err != nil {
|
||
return "", err
|
||
}
|
||
if len(infos) == 0 {
|
||
return "", fmt.Errorf("cannot get automatically,please input family_id")
|
||
}
|
||
for _, info := range infos {
|
||
if strings.Contains(y.getTokenInfo().LoginName, info.RemarkName) {
|
||
return fmt.Sprint(info.FamilyID), nil
|
||
}
|
||
}
|
||
return fmt.Sprint(infos[0].FamilyID), nil
|
||
}
|
||
|
||
// 保存家庭云中的文件到个人云
|
||
func (y *Cloud189PC) SaveFamilyFileToPersonCloud(ctx context.Context, familyId string, srcObj, dstDir model.Obj, overwrite bool) error {
|
||
// _, err := y.post(API_URL+"/family/file/saveFileToMember.action", func(req *resty.Request) {
|
||
// req.SetQueryParams(map[string]string{
|
||
// "channelId": "home",
|
||
// "familyId": familyId,
|
||
// "destParentId": destParentId,
|
||
// "fileIdList": familyFileId,
|
||
// })
|
||
// }, nil)
|
||
// return err
|
||
|
||
task := BatchTaskInfo{
|
||
FileId: srcObj.GetID(),
|
||
FileName: srcObj.GetName(),
|
||
IsFolder: BoolToNumber(srcObj.IsDir()),
|
||
}
|
||
resp, err := y.CreateBatchTask("COPY", familyId, dstDir.GetID(), map[string]string{
|
||
"groupId": "null",
|
||
"copyType": "2",
|
||
"shareId": "null",
|
||
}, task)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
|
||
for {
|
||
state, err := y.CheckBatchTask("COPY", resp.TaskID)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
switch state.TaskStatus {
|
||
case 2:
|
||
task.DealWay = IF(overwrite, 3, 2)
|
||
// 冲突时覆盖文件
|
||
if err := y.ManageBatchTask("COPY", resp.TaskID, dstDir.GetID(), task); err != nil {
|
||
return err
|
||
}
|
||
case 4:
|
||
return nil
|
||
}
|
||
time.Sleep(time.Millisecond * 400)
|
||
}
|
||
}
|
||
|
||
// 永久删除文件
|
||
func (y *Cloud189PC) Delete(ctx context.Context, familyId string, srcObj model.Obj) error {
|
||
task := BatchTaskInfo{
|
||
FileId: srcObj.GetID(),
|
||
FileName: srcObj.GetName(),
|
||
IsFolder: BoolToNumber(srcObj.IsDir()),
|
||
}
|
||
// 删除源文件
|
||
resp, err := y.CreateBatchTask("DELETE", familyId, "", nil, task)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
err = y.WaitBatchTask("DELETE", resp.TaskID, time.Second)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
// 清除回收站
|
||
resp, err = y.CreateBatchTask("CLEAR_RECYCLE", familyId, "", nil, task)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
err = y.WaitBatchTask("CLEAR_RECYCLE", resp.TaskID, time.Second)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
return nil
|
||
}
|
||
|
||
func (y *Cloud189PC) CreateBatchTask(aType string, familyID string, targetFolderId string, other map[string]string, taskInfos ...BatchTaskInfo) (*CreateBatchTaskResp, error) {
|
||
var resp CreateBatchTaskResp
|
||
_, err := y.post(API_URL+"/batch/createBatchTask.action", func(req *resty.Request) {
|
||
req.SetFormData(map[string]string{
|
||
"type": aType,
|
||
"taskInfos": MustString(utils.Json.MarshalToString(taskInfos)),
|
||
})
|
||
if targetFolderId != "" {
|
||
req.SetFormData(map[string]string{"targetFolderId": targetFolderId})
|
||
}
|
||
if familyID != "" {
|
||
req.SetFormData(map[string]string{"familyId": familyID})
|
||
}
|
||
req.SetFormData(other)
|
||
}, &resp, familyID != "")
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
return &resp, nil
|
||
}
|
||
|
||
// 检测任务状态
|
||
func (y *Cloud189PC) CheckBatchTask(aType string, taskID string) (*BatchTaskStateResp, error) {
|
||
var resp BatchTaskStateResp
|
||
_, err := y.post(API_URL+"/batch/checkBatchTask.action", func(req *resty.Request) {
|
||
req.SetFormData(map[string]string{
|
||
"type": aType,
|
||
"taskId": taskID,
|
||
})
|
||
}, &resp)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
return &resp, nil
|
||
}
|
||
|
||
// 获取冲突的任务信息
|
||
func (y *Cloud189PC) GetConflictTaskInfo(aType string, taskID string) (*BatchTaskConflictTaskInfoResp, error) {
|
||
var resp BatchTaskConflictTaskInfoResp
|
||
_, err := y.post(API_URL+"/batch/getConflictTaskInfo.action", func(req *resty.Request) {
|
||
req.SetFormData(map[string]string{
|
||
"type": aType,
|
||
"taskId": taskID,
|
||
})
|
||
}, &resp)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
return &resp, nil
|
||
}
|
||
|
||
// 处理冲突
|
||
func (y *Cloud189PC) ManageBatchTask(aType string, taskID string, targetFolderId string, taskInfos ...BatchTaskInfo) error {
|
||
_, err := y.post(API_URL+"/batch/manageBatchTask.action", func(req *resty.Request) {
|
||
req.SetFormData(map[string]string{
|
||
"targetFolderId": targetFolderId,
|
||
"type": aType,
|
||
"taskId": taskID,
|
||
"taskInfos": MustString(utils.Json.MarshalToString(taskInfos)),
|
||
})
|
||
}, nil)
|
||
return err
|
||
}
|
||
|
||
var ErrIsConflict = errors.New("there is a conflict with the target object")
|
||
|
||
// 等待任务完成
|
||
func (y *Cloud189PC) WaitBatchTask(aType string, taskID string, t time.Duration) error {
|
||
for {
|
||
state, err := y.CheckBatchTask(aType, taskID)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
switch state.TaskStatus {
|
||
case 2:
|
||
return ErrIsConflict
|
||
case 4:
|
||
return nil
|
||
}
|
||
time.Sleep(t)
|
||
}
|
||
}
|
||
|
||
func (y *Cloud189PC) getTokenInfo() *AppSessionResp {
|
||
if y.ref != nil {
|
||
return y.ref.getTokenInfo()
|
||
}
|
||
return y.tokenInfo
|
||
}
|
||
|
||
func (y *Cloud189PC) getClient() *resty.Client {
|
||
if y.ref != nil {
|
||
return y.ref.getClient()
|
||
}
|
||
return y.client
|
||
}
|
||
|
||
func (y *Cloud189PC) getCapacityInfo(ctx context.Context) (*CapacityResp, error) {
|
||
fullUrl := API_URL + "/portal/getUserSizeInfo.action"
|
||
var resp CapacityResp
|
||
_, err := y.get(fullUrl, func(req *resty.Request) {
|
||
req.SetContext(ctx)
|
||
}, &resp)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
return &resp, nil
|
||
}
|