Compare commits

..

82 Commits

Author SHA1 Message Date
MadDogOwner 874234449b fix(doubao_share): use new download info (#1890)
fix(doubao_share): update file URL retrieval to use new download info structure

Signed-off-by: MadDogOwner <xiaoran@xrgzs.top>
2026-01-01 22:29:03 +08:00
Edward 5fe267089a fix(123_open): infinite recursive call (#1854)
fix(123_open): token refresh logic

Fix token handling logic to avoid deadlock. Token method took reference of Alist's implementation.
2025-12-31 00:46:54 +08:00
KirCute 2442e302ad ci(lang): sync only new fields (#1881) 2025-12-30 15:43:30 +08:00
Tron 0612271732 fix(driver): fix file copy failure to 123pan due to incorrect etag (#1874) 2025-12-29 23:54:33 +08:00
我怎么就不是一只猫呢? c261ce78fb fix(s3): use current time as default modified time (#1860) 2025-12-29 23:52:35 +08:00
KirCute 7398e7d45e feat(alias): support load balance (#1767)
* feat(alias): support load balance

* feat(alias): support storage match for load balance

* feat(patch): add alias addition upgrade patch

* fix bugs

* fix(op/balance): optimize compatibility

* chore: change default read conflict policy

* feat(alias): refactor Alias initialization and enhance path handling

* feat(alias): enhance object masking and add support for operation restrictions

* feat(alias): enhance object masking

* feat(fs): add permission checks

* improve parsing

* update object masks

* feat(fs): enhance virtual file handling

* feat(storage): enhance virtual file retrieval and path handling

* refactor(alias): rename path handling functions for clarity and consistency

* fix(alias): update path handling in Other method to use balanced path

* fix bug

* feat(alias): add file size validation

* feat(alias): add hash consistency check

* 移除哈希合并,

* fix(alias): wrong behavior for all_strict/deterministic_or_all

* Revert "fix(alias): wrong behavior for all_strict/deterministic_or_all"

This reverts commit f001f2dcd7.

* fix(alias): wrong behavior for all_strict/deterministic_or_all

* feat(alias): support part-based read load balance

* fix(alias): list panic when leak conflict path

* fix(alias): remove Other load balance

* fix(alias): 修复 Link 方法中 resultLink 的返回类型和内容复制问题

* fix(alias): 更好的下载并发?

* chore(alias): all tips

* fix(alias): moving paths mismatch

---------

Co-authored-by: j2rong4cn <j2rong@qq.com>
Co-authored-by: ShenLin <773933146@qq.com>
2025-12-29 17:16:07 +08:00
KirCute 6e2d499ca9 refactor(bootstrap): fix OpenList-Mobile compile failed (#1857) 2025-12-24 18:46:13 +08:00
绎泽 4680ece2d9 docs(readme): add demo site (#1850)
Last Sync: 2025-12-22 12:39
2025-12-22 12:52:41 +08:00
Seven 8a4f3769d8 feat(strm): add save local mode (#1814)
* feat(strm): add KeepSameNameOnly logic

* chore(strm): skip update strm file when keepLocalDownloadFile

* feat(strm): add save local mode
2025-12-22 10:12:21 +08:00
foxxorcat cc5172e70b fix(weiyun): update sdk and support getDetails (#1845) 2025-12-22 00:20:07 +08:00
XZB-1248 a32ae97860 docs: update README for zh-CN (#1844) 2025-12-22 00:16:33 +08:00
TwoOnefour d6dd62dfe5 fix(s3): incorrect copy key with plus sign (#1820) 2025-12-22 00:15:58 +08:00
XZB-1248 216f071e64 docs: add VPS.Town as sponsor to all README (#1842)
Co-authored-by: XZB-1248 <i@1248.ink>
2025-12-21 12:26:03 +08:00
hshpy f47df5f9b2 feat(115_open): support custom pagesize (#1822) 2025-12-20 13:57:02 +08:00
KirCute ff3c4b885c fix(strm): support generate strm with sign (#1832) 2025-12-20 13:55:51 +08:00
MadDogOwner f86c7c844c feat(cloudreve_v4): add ks3 support (#1828)
Signed-off-by: MadDogOwner <xiaoran@xrgzs.top>
2025-12-19 18:05:11 +08:00
Mako (XSpy) 5db2172ed6 feat(driver): add personal / business wps drive support (#1802)
* feat(driver): add wps drive support

* feat(driver): add wps drive support

* fix(wps): update personal mode string to English

Signed-off-by: MadDogOwner <xiaoran@xrgzs.top>

* fix(wps): remove trailing slash from drive origin URL

Signed-off-by: MadDogOwner <xiaoran@xrgzs.top>

* fix(wps): correct order of options in mode selection

Signed-off-by: MadDogOwner <xiaoran@xrgzs.top>

* fix(wps): enable local sort and upload overwrite

Signed-off-by: MadDogOwner <xiaoran@xrgzs.top>

* fix(wps): resolve put bugs, fix file op problems and optimize list logic

- Fix uploading bugs. Support all uploading methods based on 8825.85d3c864.js
- Fix issues in delete/copy/move while opearting big folders.
- Use cache to optimize performance of list, especially in a deep path.

---------

Signed-off-by: MadDogOwner <xiaoran@xrgzs.top>
Co-authored-by: MadDogOwner <xiaoran@xrgzs.top>
2025-12-15 21:49:01 +08:00
wongz c4c121befc fix(139): disk-usage unmarshal failed when used capacity overflow (#1718)
Co-authored-by: Pikachu Ren <40362270+PIKACHUIM@users.noreply.github.com>
2025-12-15 21:30:08 +08:00
UcnacDx2 b4542753ba feat(drivers/139): user authentication and file batch operations (#1534)
* feat(139): Enhance 139 driver with password login and root path handling

- Added support for password-based login in the 139 driver.
- Introduced RootPath field to store the root directory path.
- Updated Init method to handle family and group types more effectively.
- Implemented new methods for handling file operations in family and group contexts.
- Enhanced error handling and logging for better debugging.
- Added new request and response structures for batch operations and document modifications.
- Improved encryption and decryption methods for secure communication.

* Update drivers/139/util.go

Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com>
Signed-off-by: UcnacDx2 <127503808+UcnacDx2@users.noreply.github.com>

* Update drivers/139/util.go

Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com>
Signed-off-by: UcnacDx2 <127503808+UcnacDx2@users.noreply.github.com>

* Update drivers/139/util.go

Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com>
Signed-off-by: UcnacDx2 <127503808+UcnacDx2@users.noreply.github.com>

* Update drivers/139/util.go

Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com>
Signed-off-by: UcnacDx2 <127503808+UcnacDx2@users.noreply.github.com>

* Update drivers/139/util.go

Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com>
Signed-off-by: UcnacDx2 <127503808+UcnacDx2@users.noreply.github.com>

---------

Signed-off-by: UcnacDx2 <127503808+UcnacDx2@users.noreply.github.com>
Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com>
2025-12-15 19:38:33 +08:00
KirCute 2a99c97d52 feat(ldap): support webdav, ftp and sftp login (#1746)
* feat(ldap): support webdav, ftp and sftp login

* fix: apply suggestions of Copilot

* feat(ldap) support ftp, sftp and webdav auto-register
2025-12-15 16:53:38 +08:00
MadDogOwner 0a407c3d8b fix(openlist): disable status check for openlist driver (#1757)
* fix(openlist): disable status check to avoid network stability issues

* fix(alist_v3): disable status check to avoid network stability issues

Signed-off-by: MadDogOwner <xiaoran@xrgzs.top>

---------

Signed-off-by: MadDogOwner <xiaoran@xrgzs.top>
2025-12-15 16:48:05 +08:00
KirCute b2596fdc24 refactor(bootstrap): move booting logic to bootstrap package (#1773)
* refactor(bootstrap): move booting to bootstrap package

* chore(log): reduce level of some callings of `utils.Log.Fatal`

* fix(s3): no shutdown after SIGTERM received

* fix: add handle hook
2025-12-15 16:47:50 +08:00
zzzhr1990 2dbe1b00d3 fix(halalcloud_open): halal-cloud upload issues (#1800)
fix halal-cloud upload issues
2025-12-15 16:47:09 +08:00
KirCute 1fc9c83df1 fix(ilanzou): parse vip size (#1792) 2025-12-12 12:13:42 +08:00
j2rong4cn d31e1a333d feat(model): add object mask support and enhance cache/task handling (#1743) 2025-12-11 15:10:42 +08:00
jenfonro e1bba7072b fix(task): tasks keep being cancelled (#1745)
* fix_cancel

* update(go.mod): update tache version

* tache v0.2.2

---------

Co-authored-by: j2rong4cn <j2rong@qq.com>
2025-12-10 19:09:16 +08:00
MadDogOwner 94c7d68413 feat(utils): add support for ignoring '@eaDir' system files (#1779) 2025-12-10 13:45:14 +08:00
KirCute 9ed77a5875 feat(driver): add AList v3 (#1721)
* feat(driver/openlist): compatible with AList v3

* Revert "feat(driver/openlist): compatible with AList v3"

This reverts commit 90f3f80186.

* feat(driver): add AList v3

* Revert "feat(patch): add migration from Alist V3 driver to OpenList (#919)"

Signed-off-by: MadDogOwner <xiaoran@xrgzs.top>

---------

Signed-off-by: MadDogOwner <xiaoran@xrgzs.top>
Co-authored-by: MadDogOwner <xiaoran@xrgzs.top>
2025-12-08 22:25:10 +08:00
varg1714 7d6d3b8f55 feat(fs): Support customizing the cache time for a specific path (#1533)
* feat(fs): Support customizing the cache time for a specific path

* feat(fs): Get the cache rule for driver information.

* feat(fs): Support globbing.

* feat(fs): Add log.

---------

Signed-off-by: ShenLin <773933146@qq.com>
Co-authored-by: ShenLin <773933146@qq.com>
2025-12-04 09:59:27 +08:00
ShenLin 5480d61f70 refactor!(userAgent): merge most userAgent into base (#1722)
refactor!(userAgent): merge all userAgent into base

1. change var to const
2. remove duplicated ua definetion after original Resty R
3. upgrade Chrome and OS versions
2025-12-04 09:54:14 +08:00
j2rong4cn 96cd714385 refactor(op): remove automatic Path assignment (#1734)
* refactor: 移除 ObjResp 中的 Id 和 Path 字段

* 移除op.List的自动设置Path
Path和Id只在驱动内使用,不应由op.List设置Path

* cnb_releases:将 Addition 结构体中的 RootPath 字段为 RootID
当List方法加载二级目录时,若使用的是Id,应对使用driver.RootID

* doubao_share: 添加潜在bug注释

* 添加 GetRootPath 方法到多个驱动
2025-12-03 00:55:40 +08:00
VXTLS e29d92f92e fix(mediafire): enable automatic session token acquisition and fix gzip parsing (#1661)
* fix(mediafire): enable automatic session token acquisition and fix gzip parsing

- Fix Init() method to allow automatic session token retrieval from cookie
- Change SessionToken from required to optional in configuration
- Add proper gzip decompression support for API responses
- Improve error handling for session token acquisition failures
- Update help text to clarify authentication requirements

Resolves initialization failure and JSON parsing errors when session token
can be automatically obtained from browser cookie.

* fix(mediafire): ensure driver files end with newline

* chore: gofmt drivers/mediafire/*.go
2025-12-02 12:13:39 +08:00
ShenLin c5f57bbcc5 fix(drivers/crypt): remove hard dependency on RemotePath (#1713) 2025-11-28 01:21:22 +08:00
j2rong4cn 9835afc645 refactor: improve upload handling (#1455)
* fix(quark): refactor upPart to use http.NewRequest

* fix(quark): improved upload handling

* fix(quark_open): improved upload handling

* fix: add retry context to multiple upload functions

* fix: optimize hash calculation in multipart upload to avoid blocking

* fix: update error handling in lifecycle functions for better clarity

* fix: update upload progress calculation to improve accuracy

* fix: simplify error handling in lifecycle functions for improved readability

* fix: remove unnecessary mutex for part uploads to simplify code

* fix(stream): simplify file handling in NewStreamSectionReader and improve error messages

* fix(terabox): optimize chunk count calculation in Put method

* perf(chaoxing): 表单上传文件0拷贝

* fix(cnb_releases): improve file upload progress tracking

* fix(baidu_netdisk): improve upload handling

* fix(upload): optimize buffer initialization for file uploads

* fix(baidu_netdisk): add retry condition to skip ErrUploadIDExpired in upload loop
2025-11-27 19:34:03 +08:00
Seven 1f373eac8d chore(strm): avoid generating empty folders (#1720)
chore(strm): empty folders are not generated locally
2025-11-27 14:02:11 +08:00
jenfonro ede96a314c fix(onedrive_shareurl): Reduce temporary file errors (#1686)
* fix onedrive_shareurl

* .
2025-11-25 21:43:17 +08:00
Seven 72206ac9f6 feat(strm): keep local download file (#1707) 2025-11-25 18:08:48 +08:00
ShenLin 62dedb2a2e fix(pkg/aria2): use pointer receivers for Call methods (#1706) 2025-11-25 13:25:42 +08:00
ShenLin 7189c5b461 chore(pkg/aria2): simplify context cancellation handling in RPC calls (#1705) 2025-11-25 13:09:04 +08:00
ShenLin 1a445f9d3f chore(archive): fix struct literal uses unkeyed fields (#1704) 2025-11-25 13:05:36 +08:00
ShenLin aa22884079 fix(search): fix duplicated variable init (#1703) 2025-11-25 12:23:25 +08:00
ImoutoHeaven 316d4caf37 feat(search): Add task queue for Meilisearch to prevent race conditions (#1423)
* Add task queue for Meilisearch to prevent race conditions

- Implement TaskQueueManager for async index operations
- Queue update tasks and process them in batches every 30 seconds
- Check pending task status before executing new operations
- Optimize batch indexing and deletion logic
- Fix type assertion bug in buildSearchDocumentFromResults

* fix(search): re-enqueue skipped tasks to prevent task loss

When tasks are skipped due to pending dependencies, they are now
re-enqueued if not already in queue. This prevents task loss while
avoiding overwriting newer snapshots for the same parent.

* fix(copilot-comment): Invoke Stop() & err of SliceConvert

---------

Co-authored-by: ImoutoHeaven <noreply@imoutoheaven.org>
Co-authored-by: jyxjjj <773933146@qq.com>
2025-11-25 11:38:27 +08:00
KirCute 60a489eb68 feat(archive): support non-overwrite decompress (#1701) 2025-11-25 10:27:29 +08:00
varg1714 b22e211044 feat(fs): Add skipExisting option to move and copy, merge option to copy (#1556)
* fix(fs): Add skipExisting option to move and copy.

* feat(fs): Add merge option to copy.

* feat(fs): Code smell.

* feat(fs): Code smell.
2025-11-24 14:20:24 +08:00
KirCute ca401b9af9 fix(local): assign non-CoW copy requests to the task module (#1669)
* fix(local): assign non-CoW copy requests to the task module

* fix build

* fix cross device
2025-11-24 14:14:53 +08:00
jenfonro addce8b691 feat(baidu_netdisk): Add shard upload timeout setting (#1682)
add timeout
2025-11-24 14:14:34 +08:00
Seven 42fc841dc1 feat(strm): custom path prefixes (#1697)
fix(strm): custom path prefixes

Signed-off-by: ShenLin <773933146@qq.com>
Co-authored-by: ShenLin <773933146@qq.com>
2025-11-24 14:11:59 +08:00
varg1714 4c0916b64b fix(strm): fix the name and type issue (#1630)
* fix(strm): fix the name and type issue

* fix(strm): update version
2025-11-24 14:05:49 +08:00
VXTLS 3989d35abd fix(misskey): folderId format validation and root directory handling (#1647)
fix(misskey): Fix folderId format validation and root directory handling
2025-11-21 12:18:54 +08:00
KirCute 72e2ae1f14 feat(fs): support manually trigger objs update hook (#1620)
* feat(fs): support manually trigger objs update hook

* fix: support driver internal copy & move case

* fix

* fix: apply suggestions of Copilot
2025-11-21 12:18:20 +08:00
Seven 3e37f575d8 fix(openlist_driver): ensure UA is correctly propagated (#1679) 2025-11-21 12:13:41 +08:00
MoYan c0d480366d fix(driver/123): initialize Platform field (#1644)
* fix(driver/123): initialization the Platform field

Signed-off-by: MoYan <1561515308@qq.com>

* Fix formatting of Platform field in Pan123

Signed-off-by: MoYan <1561515308@qq.com>

---------

Signed-off-by: MoYan <1561515308@qq.com>
2025-11-14 18:49:44 +08:00
Copilot 9de7561154 feat(upload): add optional system file filtering for uploads (#1634) 2025-11-14 14:45:39 +08:00
MadDogOwner 0866b9075f fix(link): correct link cache mode bitwise comparison (#1635)
* fix(link): correct link cache mode bitwise comparison

Signed-off-by: MadDogOwner <xiaoran@xrgzs.top>

* refactor(link): use explicit flag equality for link cache mode bitmask checks

Signed-off-by: MadDogOwner <xiaoran@xrgzs.top>

---------

Signed-off-by: MadDogOwner <xiaoran@xrgzs.top>
2025-11-13 13:52:33 +08:00
KirCute 055696f576 feat(s3): support frontend direct upload (#1631)
* feat(s3): support frontend direct upload

* feat(s3): support custom direct upload host

* fix: apply suggestions of Copilot
2025-11-13 13:22:17 +08:00
ShenLin 854415160c chore(issue templates): require logs (#1626) 2025-11-12 13:04:13 +08:00
varg1714 8f4f7d1291 feat(doubao): Add rate limiting (#1618) 2025-11-11 21:59:10 +08:00
KirCute ee2c77acd8 fix(archive/zip): user specific encoding for non-EFS zips (#1599)
* fix(archive/zip): user specific encoding for non-EFS zips

* fix(stream): simplify head cache initialization and improve reader retrieval logic

* fix: support multipart zips (.z01)

* chore(deps): update github.com/KirCute/zip to v1.0.1

---------

Co-authored-by: j2rong4cn <j2rong@qq.com>
Co-authored-by: Pikachu Ren <40362270+PIKACHUIM@users.noreply.github.com>
2025-11-10 19:08:50 +08:00
yuyamionini fc90ec1b53 fix(terabox): wrong return code used (#1547)
fix(terabox): rename, delete, copy operations sometimes failed

Signed-off-by: yuyamionini <46483865+yuyamionini@users.noreply.github.com>
2025-11-10 13:40:00 +08:00
jenfonro 7d78944d14 fix(baidu_netdisk): Fix Baidu Netdisk resume uploads sticking to the same upload host (#1609)
Fix Baidu Netdisk resume uploads sticking to the same upload host
2025-11-09 20:43:02 +08:00
jenfonro f2e0fe8589 refactor(fs): implement immediate retry within task execution cycle (#1575) 2025-11-07 19:11:11 +08:00
ASLant 39dcf9bd19 feat(onedrive): support frontend direct upload (#1532)
* OneDrive添加直连上传

* refactor

* fix: duplicate root path join

---------

Co-authored-by: KirCute <951206789@qq.com>
2025-11-06 23:22:02 +08:00
Seven 25f38df4ca fix(strm): non-specified type generates strm (#1585)
* fix(strm): non-specified type generates strm

* fix(strm): only insert to strmTrie if SaveStrmToLocal is enabled

* fix(strm): update suffix handling in convert2strmObjs function

* fix(strm): refactor generateStrm to use range reader

---------

Co-authored-by: j2rong4cn <j2rong@qq.com>
2025-11-06 20:58:43 +08:00
j2rong4cn a1f1f98f94 refactor(stream): simplify code (#1590)
* refactor(stream): simplify Close method and update SeekableStream to use RangeReader interface

* refactor(stream):  improve RangeRead comments for clarity
2025-11-06 20:06:48 +08:00
KirCute affc499913 fix(189): disk-usage unmarshal failed when used capacity overflow (#1577) 2025-11-05 12:35:51 +08:00
KirCute c7574b545c feat(github_release): support Source code (zip/tar.gz) (#1581)
* support Github Release Source code (zip/tar.gz)

* fix TarballUrl and ZipballUrl

* fix show source code by allversion

---------

Co-authored-by: nibazshab <44338441+nibazshab@users.noreply.github.com>
2025-11-05 12:30:05 +08:00
hcrgm 9e852ba12d fix(baidu_netdisk): improve upload experience (#1562)
* fix(baidu_netdisk): improve upload experience

* fix(typo): URL should be uppercase, apply suggestion from @Copilot

Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com>
Signed-off-by: ShenLin <773933146@qq.com>

* fix(typo): URL should be uppercase, apply suggestion from @Copilot

Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com>
Signed-off-by: ShenLin <773933146@qq.com>

* fix(baidu_netdisk): use "UploadAPI" as a fallback when using dynamic upload api

* fix(baidu_netdisk): all uploads share the same upload url cache

* fix(drivers/baidu_netdisk): defer uploadUrlMu unlock

* update driver.go to main

---------

Signed-off-by: ShenLin <773933146@qq.com>
Signed-off-by: jenfonro <799170122@qq.com>
Co-authored-by: ShenLin <773933146@qq.com>
Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com>
Co-authored-by: jenfonro <799170122@qq.com>
2025-11-05 12:21:32 +08:00
j2rong4cn 174eae802a perf(stream): optimize CacheFullAndWriter for better memory management (#1584)
* perf(stream): optimize CacheFullAndWriter for better memory management

* fix(stream): ensure proper seek handling in CacheFullAndWriter for improved data integrity
2025-11-05 12:16:09 +08:00
KirCute b9f058fcc9 fix(backup-restore): add shares (#1500) 2025-11-05 12:11:20 +08:00
j2rong4cn 6de15b6310 feat(stream): enhance GetRangeReaderFromLink rate limiting (#1528)
* feat(stream): enhance GetRangeReaderFromLink rate limiting

* refactor(stream): update GetRangeReaderFromMFile to return *model.FileRangeReader

* refactor(stream): simplify context error handling in RateLimitReader, RateLimitWriter, and RateLimitFile

* refactor(net): replace custom LimitedReadCloser with readers.NewLimitedReadCloser

* fix(model): update Link.ContentLength JSON tag for correct serialization

* docs(model): add clarification to FileRangeReader usage comment
2025-11-04 23:56:09 +08:00
jenfonro 2844797684 fix(baidu_netdisk): support resuming uploads when an error occurs (#1279)
support resuming uploads when an error occurs
2025-11-04 13:33:33 +08:00
Seven 9f4e439478 chore(strm): Built-in file types support modification (#1483) 2025-11-04 10:33:16 +08:00
jenfonro 9d09ee133d fix(google_driver): fix google link file display size (#1335)
* fix file link display size

* fix performance and field

* cn to en notes

---------

Co-authored-by: ShenLin <773933146@qq.com>
2025-11-04 09:26:20 +08:00
jenfonro d88f0e8f3c feat(net): support proxy configuration via config file (#1359)
* support proxy

* debug

* debug2

* del debug

* add proxy configuration with env var fallback

* comments to en

* refactor(env): fallback env

---------

Co-authored-by: jyxjjj <773933146@qq.com>
2025-11-04 09:01:35 +08:00
ex-hentai 0857478516 feat(thunder): allow setting space (#1219)
allows access to files on remote devices via Thunder's tunneling service.
2025-11-03 10:53:38 +08:00
Seven 66d9809057 feat(strm): strm local file (#1127)
* feat(strm): strm local file

* feat: 代码优化

* feat: 访问被strm挂载路径时也更新

* fix: 路径最后带/判断缺失

* fix: 路径最后带/判断缺失

* refactor

* refactor

* fix: close seekable-stream in `generateStrm`

* refactor: lazy create local file

* 优化路径判断

---------

Co-authored-by: KirCute <kircute@foxmail.com>
2025-11-03 10:48:15 +08:00
MoYan db8a7e8caf feat(123): allow modification of the platform header (#1542)
* feat(drivers/123): Allow modification of the platform field

* feat(drivers/123): Set login platfrom as web

* fix(drivers/123): update platform field help value
2025-11-03 10:02:52 +08:00
walloo 8f18e34da0 fix(alias): nil panic in ResolveLinkCacheMode (#1527)
* fix(alias): Check the driver path during initialization

* fix(alias): Don't check the driver path during initialization anymore.
2025-10-23 22:18:56 +08:00
ShenLin 525f26dc23 feat(command): add --config flag to set custom config path (#1479) 2025-10-22 19:35:34 +08:00
NewbieOrange a0fcfa3ed2 fix(aliyundrive_open): use safe disk usage calculation (#1510) 2025-10-20 22:05:14 +08:00
ILoveScratch 15f276537c fix(share): remove share when user delete (#1493) 2025-10-19 22:48:21 +08:00
MadDogOwner 623a12050e feat(openlist): add PassIPToUpsteam to driver (#1498) 2025-10-19 22:45:40 +08:00
243 changed files with 9071 additions and 3123 deletions
+9 -7
View File
@@ -13,7 +13,7 @@ body:
attributes:
label: 请确认以下事项
description: |
您必须勾选以下内容,否则您的问题可能会被直接关闭。
您必须确认、同意并勾选以下内容,否则您的问题一定会被直接关闭。
或者您可以去[讨论区](https://github.com/OpenListTeam/OpenList/discussions)。
options:
- label: |
@@ -59,6 +59,14 @@ body:
label: 问题描述(必填)
validations:
required: true
- type: textarea
id: logs
attributes:
label: 日志(必填)
description: |
请复制粘贴错误日志,或者截图。(可隐藏隐私字段) [查看方法](https://doc.oplist.org/faq/howto#%E5%A6%82%E4%BD%95%E5%BF%AB%E9%80%9F%E5%AE%9A%E4%BD%8Dbug)
validations:
required: true
- type: textarea
id: config
attributes:
@@ -67,12 +75,6 @@ body:
请提供您的`OpenList`应用的配置文件,并截图相关存储配置。(可隐藏隐私字段)
validations:
required: true
- type: textarea
id: logs
attributes:
label: 日志(可选)
description: |
请复制粘贴错误日志,或者截图。(可隐藏隐私字段) [查看方法](https://doc.oplist.org/faq/howto#%E5%A6%82%E4%BD%95%E5%BF%AB%E9%80%9F%E5%AE%9A%E4%BD%8Dbug)
- type: textarea
id: reproduction
attributes:
+9 -7
View File
@@ -13,7 +13,7 @@ body:
attributes:
label: Please confirm the following
description: |
You must check all the following, otherwise your issue may be closed directly.
You must confirm, agree, and check all the following, otherwise your issue will definitely be closed directly.
Or you can go to the [discussions](https://github.com/OpenListTeam/OpenList/discussions).
options:
- label: |
@@ -59,6 +59,14 @@ body:
label: Bug Description (required)
validations:
required: true
- type: textarea
id: logs
attributes:
label: Logs (required)
description: |
Please copy and paste any relevant log output or screenshots. (You may mask sensitive fields) [Guide](https://doc.oplist.org/faq/howto#how-to-quickly-locate-bugs)
validations:
required: true
- type: textarea
id: config
attributes:
@@ -67,12 +75,6 @@ body:
Please provide your `OpenList` application's configuration file and a screenshot of the relevant storage configuration. (You may mask sensitive fields)
validations:
required: true
- type: textarea
id: logs
attributes:
label: Logs (optional)
description: |
Please copy and paste any relevant log output or screenshots. (You may mask sensitive fields) [Guide](https://doc.oplist.org/faq/howto#how-to-quickly-locate-bugs)
- type: textarea
id: reproduction
attributes:
+6 -1
View File
@@ -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
[![VPS.Town](https://vps.town/static/images/sponsor.png)](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
View File
@@ -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* 仅用于错误报告和功能请求。**
## 赞助者
[![VPS.Town](https://vps.town/static/images/sponsor.png)](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
View File
@@ -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* はバグ報告と機能リクエスト専用です。**
## スポンサー
[![VPS.Town](https://vps.town/static/images/sponsor.png)](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
View File
@@ -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
[![VPS.Town](https://vps.town/static/images/sponsor.png)](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
View File
@@ -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
View File
@@ -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
View File
@@ -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
View File
@@ -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
+1
View File
@@ -2,6 +2,7 @@ package flags
var (
DataDir string
ConfigPath string
Debug bool
NoPrefix bool
Dev bool
+20 -3
View File
@@ -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{}])
+2 -1
View File
@@ -27,7 +27,8 @@ func Execute() {
}
func init() {
RootCmd.PersistentFlags().StringVar(&flags.DataDir, "data", "data", "data folder")
RootCmd.PersistentFlags().StringVar(&flags.DataDir, "data", "data", "data directory (relative paths are resolved against the current working directory)")
RootCmd.PersistentFlags().StringVar(&flags.ConfigPath, "config", "", "path to config.json (relative to current working directory; defaults to [data directory]/config.json, where [data directory] is set by --data)")
RootCmd.PersistentFlags().BoolVar(&flags.Debug, "debug", false, "start with debug mode")
RootCmd.PersistentFlags().BoolVar(&flags.NoPrefix, "no-prefix", false, "disable env prefix")
RootCmd.PersistentFlags().BoolVar(&flags.Dev, "dev", false, "start with dev mode")
+4 -239
View File
@@ -1,34 +1,13 @@
package cmd
import (
"context"
"errors"
"fmt"
"net"
"net/http"
"os"
"os/signal"
"strconv"
"sync"
"syscall"
"time"
"github.com/OpenListTeam/OpenList/v4/cmd/flags"
"github.com/OpenListTeam/OpenList/v4/internal/bootstrap"
"github.com/OpenListTeam/OpenList/v4/internal/conf"
"github.com/OpenListTeam/OpenList/v4/internal/fs"
"github.com/OpenListTeam/OpenList/v4/pkg/utils"
"github.com/OpenListTeam/OpenList/v4/server"
"github.com/OpenListTeam/OpenList/v4/server/middlewares"
"github.com/OpenListTeam/sftpd-openlist"
ftpserver "github.com/fclairamb/ftpserverlib"
"github.com/gin-gonic/gin"
log "github.com/sirupsen/logrus"
"github.com/spf13/cobra"
"golang.org/x/net/http2"
"golang.org/x/net/http2/h2c"
"github.com/quic-go/quic-go/http3"
)
// ServerCmd represents the server command
@@ -38,161 +17,9 @@ var ServerCmd = &cobra.Command{
Long: `Start the server at the specified address
the address is defined in config file`,
Run: func(cmd *cobra.Command, args []string) {
Init()
if conf.Conf.DelayedStart != 0 {
utils.Log.Infof("delayed start for %d seconds", conf.Conf.DelayedStart)
time.Sleep(time.Duration(conf.Conf.DelayedStart) * time.Second)
}
bootstrap.InitOfflineDownloadTools()
bootstrap.LoadStorages()
bootstrap.InitTaskManager()
if !flags.Debug && !flags.Dev {
gin.SetMode(gin.ReleaseMode)
}
r := gin.New()
// gin log
if conf.Conf.Log.Filter.Enable {
r.Use(middlewares.FilteredLogger())
} else {
r.Use(gin.LoggerWithWriter(log.StandardLogger().Out))
}
r.Use(gin.RecoveryWithWriter(log.StandardLogger().Out))
server.Init(r)
var httpHandler http.Handler = r
if conf.Conf.Scheme.EnableH2c {
httpHandler = h2c.NewHandler(r, &http2.Server{})
}
var httpSrv, httpsSrv, unixSrv *http.Server
var quicSrv *http3.Server
if conf.Conf.Scheme.HttpPort != -1 {
httpBase := fmt.Sprintf("%s:%d", conf.Conf.Scheme.Address, conf.Conf.Scheme.HttpPort)
fmt.Printf("start HTTP server @ %s\n", httpBase)
utils.Log.Infof("start HTTP server @ %s", httpBase)
httpSrv = &http.Server{Addr: httpBase, Handler: httpHandler}
go func() {
err := httpSrv.ListenAndServe()
if err != nil && !errors.Is(err, http.ErrServerClosed) {
utils.Log.Fatalf("failed to start http: %s", err.Error())
}
}()
}
if conf.Conf.Scheme.HttpsPort != -1 {
httpsBase := fmt.Sprintf("%s:%d", conf.Conf.Scheme.Address, conf.Conf.Scheme.HttpsPort)
fmt.Printf("start HTTPS server @ %s\n", httpsBase)
utils.Log.Infof("start HTTPS server @ %s", httpsBase)
httpsSrv = &http.Server{Addr: httpsBase, Handler: r}
go func() {
err := httpsSrv.ListenAndServeTLS(conf.Conf.Scheme.CertFile, conf.Conf.Scheme.KeyFile)
if err != nil && !errors.Is(err, http.ErrServerClosed) {
utils.Log.Fatalf("failed to start https: %s", err.Error())
}
}()
if conf.Conf.Scheme.EnableH3 {
fmt.Printf("start HTTP3 (quic) server @ %s\n", httpsBase)
utils.Log.Infof("start HTTP3 (quic) server @ %s", httpsBase)
r.Use(func(c *gin.Context) {
if c.Request.TLS != nil {
port := conf.Conf.Scheme.HttpsPort
c.Header("Alt-Svc", fmt.Sprintf("h3=\":%d\"; ma=86400", port))
}
c.Next()
})
quicSrv = &http3.Server{Addr: httpsBase, Handler: r}
go func() {
err := quicSrv.ListenAndServeTLS(conf.Conf.Scheme.CertFile, conf.Conf.Scheme.KeyFile)
if err != nil && !errors.Is(err, http.ErrServerClosed) {
utils.Log.Fatalf("failed to start http3 (quic): %s", err.Error())
}
}()
}
}
if conf.Conf.Scheme.UnixFile != "" {
fmt.Printf("start unix server @ %s\n", conf.Conf.Scheme.UnixFile)
utils.Log.Infof("start unix server @ %s", conf.Conf.Scheme.UnixFile)
unixSrv = &http.Server{Handler: httpHandler}
go func() {
listener, err := net.Listen("unix", conf.Conf.Scheme.UnixFile)
if err != nil {
utils.Log.Fatalf("failed to listen unix: %+v", err)
}
// set socket file permission
mode, err := strconv.ParseUint(conf.Conf.Scheme.UnixFilePerm, 8, 32)
if err != nil {
utils.Log.Errorf("failed to parse socket file permission: %+v", err)
} else {
err = os.Chmod(conf.Conf.Scheme.UnixFile, os.FileMode(mode))
if err != nil {
utils.Log.Errorf("failed to chmod socket file: %+v", err)
}
}
err = unixSrv.Serve(listener)
if err != nil && !errors.Is(err, http.ErrServerClosed) {
utils.Log.Fatalf("failed to start unix: %s", err.Error())
}
}()
}
if conf.Conf.S3.Port != -1 && conf.Conf.S3.Enable {
s3r := gin.New()
s3r.Use(gin.LoggerWithWriter(log.StandardLogger().Out), gin.RecoveryWithWriter(log.StandardLogger().Out))
server.InitS3(s3r)
s3Base := fmt.Sprintf("%s:%d", conf.Conf.Scheme.Address, conf.Conf.S3.Port)
fmt.Printf("start S3 server @ %s\n", s3Base)
utils.Log.Infof("start S3 server @ %s", s3Base)
go func() {
var err error
if conf.Conf.S3.SSL {
httpsSrv = &http.Server{Addr: s3Base, Handler: s3r}
err = httpsSrv.ListenAndServeTLS(conf.Conf.Scheme.CertFile, conf.Conf.Scheme.KeyFile)
}
if !conf.Conf.S3.SSL {
httpSrv = &http.Server{Addr: s3Base, Handler: s3r}
err = httpSrv.ListenAndServe()
}
if err != nil && !errors.Is(err, http.ErrServerClosed) {
utils.Log.Fatalf("failed to start s3 server: %s", err.Error())
}
}()
}
var ftpDriver *server.FtpMainDriver
var ftpServer *ftpserver.FtpServer
if conf.Conf.FTP.Listen != "" && conf.Conf.FTP.Enable {
var err error
ftpDriver, err = server.NewMainDriver()
if err != nil {
utils.Log.Fatalf("failed to start ftp driver: %s", err.Error())
} else {
fmt.Printf("start ftp server on %s\n", conf.Conf.FTP.Listen)
utils.Log.Infof("start ftp server on %s", conf.Conf.FTP.Listen)
go func() {
ftpServer = ftpserver.NewFtpServer(ftpDriver)
err = ftpServer.ListenAndServe()
if err != nil {
utils.Log.Fatalf("problem ftp server listening: %s", err.Error())
}
}()
}
}
var sftpDriver *server.SftpDriver
var sftpServer *sftpd.SftpServer
if conf.Conf.SFTP.Listen != "" && conf.Conf.SFTP.Enable {
var err error
sftpDriver, err = server.NewSftpDriver()
if err != nil {
utils.Log.Fatalf("failed to start sftp driver: %s", err.Error())
} else {
fmt.Printf("start sftp server on %s", conf.Conf.SFTP.Listen)
utils.Log.Infof("start sftp server on %s", conf.Conf.SFTP.Listen)
go func() {
sftpServer = sftpd.NewSftpServer(sftpDriver)
err = sftpServer.RunServer()
if err != nil {
utils.Log.Fatalf("problem sftp server listening: %s", err.Error())
}
}()
}
}
bootstrap.Init()
defer bootstrap.Release()
bootstrap.Start()
// Wait for interrupt signal to gracefully shutdown the server with
// a timeout of 1 second.
quit := make(chan os.Signal, 1)
@@ -201,69 +28,7 @@ the address is defined in config file`,
// kill -9 is syscall. SIGKILL but can"t be catch, so don't need add it
signal.Notify(quit, syscall.SIGINT, syscall.SIGTERM)
<-quit
utils.Log.Println("Shutdown server...")
fs.ArchiveContentUploadTaskManager.RemoveAll()
Release()
ctx, cancel := context.WithTimeout(context.Background(), 1*time.Second)
defer cancel()
var wg sync.WaitGroup
if conf.Conf.Scheme.HttpPort != -1 {
wg.Add(1)
go func() {
defer wg.Done()
if err := httpSrv.Shutdown(ctx); err != nil {
utils.Log.Fatal("HTTP server shutdown err: ", err)
}
}()
}
if conf.Conf.Scheme.HttpsPort != -1 {
wg.Add(1)
go func() {
defer wg.Done()
if err := httpsSrv.Shutdown(ctx); err != nil {
utils.Log.Fatal("HTTPS server shutdown err: ", err)
}
}()
if conf.Conf.Scheme.EnableH3 {
wg.Add(1)
go func() {
defer wg.Done()
if err := quicSrv.Shutdown(ctx); err != nil {
utils.Log.Fatal("HTTP3 (quic) server shutdown err: ", err)
}
}()
}
}
if conf.Conf.Scheme.UnixFile != "" {
wg.Add(1)
go func() {
defer wg.Done()
if err := unixSrv.Shutdown(ctx); err != nil {
utils.Log.Fatal("Unix server shutdown err: ", err)
}
}()
}
if conf.Conf.FTP.Listen != "" && conf.Conf.FTP.Enable && ftpServer != nil && ftpDriver != nil {
wg.Add(1)
go func() {
defer wg.Done()
ftpDriver.Stop()
if err := ftpServer.Stop(); err != nil {
utils.Log.Fatal("FTP server shutdown err: ", err)
}
}()
}
if conf.Conf.SFTP.Listen != "" && conf.Conf.SFTP.Enable && sftpServer != nil && sftpDriver != nil {
wg.Add(1)
go func() {
defer wg.Done()
if err := sftpServer.Close(); err != nil {
utils.Log.Fatal("SFTP server shutdown err: ", err)
}
}()
}
wg.Wait()
utils.Log.Println("Server exit")
bootstrap.Shutdown(1 * time.Second)
},
}
+7 -6
View File
@@ -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)
+1 -1
View File
@@ -17,7 +17,7 @@ type Addition struct {
var config = driver.Config{
Name: "115 Cloud",
DefaultRoot: "0",
LinkCacheType: 2,
LinkCacheMode: driver.LinkCacheUA,
}
func init() {
+7 -1
View File
@@ -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 {
+2 -1
View File
@@ -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"`
}
@@ -19,7 +20,7 @@ type Addition struct {
var config = driver.Config{
Name: "115 Open",
DefaultRoot: "0",
LinkCacheType: 2,
LinkCacheMode: driver.LinkCacheUA,
}
func init() {
+2 -2
View File
@@ -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))
+3 -1
View File
@@ -41,7 +41,9 @@ func (d *Pan123) GetAddition() driver.Additional {
}
func (d *Pan123) Init(ctx context.Context) error {
_, err := d.Request(UserInfo, http.MethodGet, nil, nil)
_, err := d.Request(UserInfo, http.MethodGet, func(req *resty.Request) {
req.SetHeader("platform", "web")
}, nil)
return err
}
+3 -1
View File
@@ -12,7 +12,8 @@ type Addition struct {
//OrderBy string `json:"order_by" type:"select" options:"file_id,file_name,size,update_at" default:"file_name"`
//OrderDirection string `json:"order_direction" type:"select" options:"asc,desc" default:"asc"`
AccessToken string
UploadThread int `json:"UploadThread" type:"number" default:"3" help:"the threads of upload"`
UploadThread int `json:"UploadThread" type:"number" default:"3" help:"the threads of upload"`
Platform string `json:"platform" type:"string" default:"web" help:"the platform header value, sent with API requests"`
}
var config = driver.Config{
@@ -27,6 +28,7 @@ func init() {
return &Pan123{
Addition: Addition{
UploadThread: 3,
Platform: "web",
},
}
})
+7 -16
View File
@@ -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
},
+1 -1
View File
@@ -203,7 +203,7 @@ do:
"referer": "https://www.123pan.com/",
"authorization": "Bearer " + d.AccessToken,
"user-agent": "Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) openlist-client",
"platform": "web",
"platform": d.Platform,
"app-version": "3",
//"user-agent": base.UserAgent,
})
+4
View File
@@ -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)
+19
View File
@@ -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
}
+1 -1
View File
@@ -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可以兼容
+115
View File
@@ -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
View File
@@ -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
},
+8 -44
View File
@@ -13,7 +13,6 @@ import (
"time"
"github.com/OpenListTeam/OpenList/v4/drivers/base"
"github.com/OpenListTeam/OpenList/v4/internal/op"
"github.com/go-resty/resty/v2"
"github.com/google/uuid"
log "github.com/sirupsen/logrus"
@@ -22,8 +21,6 @@ import (
var ( // 不同情况下获取的AccessTokenQPS限制不同 如下模块化易于拓展
Api = "https://open-api.123pan.com"
AccessToken = InitApiInfo(Api+"/api/v1/access_token", 1)
RefreshToken = InitApiInfo(Api+"/api/v1/oauth2/access_token", 1)
UserInfo = InitApiInfo(Api+"/api/v1/user/info", 1)
FileList = InitApiInfo(Api+"/api/v2/file/list", 3)
DownloadInfo = InitApiInfo(Api+"/api/v1/file/download_info", 5)
@@ -40,11 +37,14 @@ var ( // 不同情况下获取的AccessTokenQPS限制不同 如下模块化易
)
func (d *Open123) Request(apiInfo *ApiInfo, method string, callback base.ReqCallback, resp interface{}) ([]byte, error) {
retryToken := true
for {
token, err := d.getAccessToken(false)
if err != nil {
return nil, err
}
req := base.RestyClient.R()
req.SetHeaders(map[string]string{
"authorization": "Bearer " + d.AccessToken,
"authorization": "Bearer " + token,
"platform": "open_platform",
"Content-Type": "application/json",
})
@@ -74,9 +74,9 @@ func (d *Open123) Request(apiInfo *ApiInfo, method string, callback base.ReqCall
if baseResp.Code == 0 {
return body, nil
} else if baseResp.Code == 401 && retryToken {
retryToken = false
if err := d.flushAccessToken(); err != nil {
} else if baseResp.Code == 401 {
// 强制刷新Token, 有小概率会 race condition 导致多次刷新Token,但不影响正确运行
if _, err := d.getAccessToken(true); err != nil {
return nil, err
}
} else if baseResp.Code == 429 {
@@ -88,42 +88,6 @@ func (d *Open123) Request(apiInfo *ApiInfo, method string, callback base.ReqCall
}
}
func (d *Open123) flushAccessToken() error {
if d.ClientID != "" {
if d.RefreshToken != "" {
var resp RefreshTokenResp
_, err := d.Request(RefreshToken, http.MethodPost, func(req *resty.Request) {
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)
}, &resp)
if err != nil {
return err
}
d.AccessToken = resp.AccessToken
d.RefreshToken = resp.RefreshToken
op.MustSaveDriverStorage(d)
} else if d.ClientSecret != "" {
var resp AccessTokenResp
_, err := d.Request(AccessToken, http.MethodPost, func(req *resty.Request) {
req.SetBody(base.Json{
"clientID": d.ClientID,
"clientSecret": d.ClientSecret,
})
}, &resp)
if err != nil {
return err
}
d.AccessToken = resp.Data.AccessToken
op.MustSaveDriverStorage(d)
}
}
return nil
}
func (d *Open123) SignURL(originURL, privateKey string, uid uint64, validDuration time.Duration) (newURL string, err error) {
// 生成Unix时间戳
ts := time.Now().Add(validDuration).Unix()
+125 -23
View File
@@ -14,6 +14,7 @@ import (
"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"
streamPkg "github.com/OpenListTeam/OpenList/v4/internal/stream"
"github.com/OpenListTeam/OpenList/v4/pkg/cron"
"github.com/OpenListTeam/OpenList/v4/pkg/utils"
@@ -28,6 +29,7 @@ type Yun139 struct {
Account string
ref *Yun139
PersonalCloudHost string
RootPath string
}
func (d *Yun139) Config() driver.Config {
@@ -41,7 +43,16 @@ func (d *Yun139) GetAddition() driver.Additional {
func (d *Yun139) Init(ctx context.Context) error {
if d.ref == nil {
if len(d.Authorization) == 0 {
return fmt.Errorf("authorization is empty")
if d.Username != "" && d.Password != "" {
log.Infof("139yun: authorization is empty, trying to login with password.")
newAuth, err := d.loginWithPassword()
log.Debugf("newAuth: Ok: %s", newAuth)
if err != nil {
return fmt.Errorf("login with password failed: %w", err)
}
} else {
return fmt.Errorf("authorization is empty and username/password is not provided")
}
}
err := d.refreshToken()
if err != nil {
@@ -92,7 +103,22 @@ func (d *Yun139) Init(ctx context.Context) error {
if len(d.Addition.RootFolderID) == 0 {
d.RootFolderID = d.CloudID
}
_, err := d.groupGetFiles(d.RootFolderID)
if err != nil {
return err
}
case MetaFamily:
if len(d.Addition.RootFolderID) == 0 {
// Attempt to obtain data.path as the root via a query and persist it.
if root, err := d.getFamilyRootPath(d.CloudID); err == nil && root != "" {
d.RootFolderID = root
op.MustSaveDriverStorage(d)
}
}
_, err := d.familyGetFiles(d.RootFolderID)
if err != nil {
return err
}
default:
return errs.NotImplement
}
@@ -279,6 +305,42 @@ func (d *Yun139) Move(ctx context.Context, srcObj, dstDir model.Obj) (model.Obj,
return nil, err
}
return srcObj, nil
case MetaFamily:
pathname := "/isbo/openApi/createBatchOprTask"
var contentList []string
var catalogList []string
if srcObj.IsDir() {
catalogList = append(catalogList, path.Join(srcObj.GetPath(), srcObj.GetID()))
} else {
contentList = append(contentList, path.Join(srcObj.GetPath(), srcObj.GetID()))
}
body := base.Json{
"catalogList": catalogList,
"accountInfo": base.Json{
"accountName": d.getAccount(),
"accountType": "1",
},
"contentList": contentList,
"destCatalogID": dstDir.GetID(),
"destGroupID": d.CloudID,
"destPath": path.Join(dstDir.GetPath(), dstDir.GetID()),
"destType": 0,
"srcGroupID": d.CloudID,
"srcType": 0,
"taskType": 3,
}
var resp CreateBatchOprTaskResp
_, err := d.isboPost(pathname, body, &resp)
if err != nil {
return nil, err
}
log.Debugf("[139] Move MetaFamily CreateBatchOprTaskResp.Result.ResultCode: %s", resp.Result.ResultCode)
if resp.Result.ResultCode != "0" {
return nil, fmt.Errorf("failed to move in family cloud: %s", resp.Result.ResultDesc)
}
return srcObj, nil
default:
return nil, errs.NotImplement
}
@@ -353,19 +415,27 @@ func (d *Yun139) Rename(ctx context.Context, srcObj model.Obj, newName string) e
var data base.Json
var pathname string
if srcObj.IsDir() {
// 网页接口不支持重命名家庭云文件夹
// data = base.Json{
// "catalogType": 3,
// "catalogID": srcObj.GetID(),
// "catalogName": newName,
// "commonAccountInfo": base.Json{
// "account": d.getAccount(),
// "accountType": 1,
// },
// "path": srcObj.GetPath(),
// }
// pathname = "/orchestration/familyCloud-rebuild/photoContent/v1.0/modifyCatalogInfo"
return errs.NotImplement
pathname = "/modifyCloudDocV2"
data = base.Json{
"catalogType": 3,
"cloudID": d.CloudID,
"commonAccountInfo": base.Json{
"account": d.getAccount(),
"accountType": "1",
},
"docLibName": newName,
"docLibraryID": srcObj.GetID(),
"path": path.Join(srcObj.GetPath(), srcObj.GetID()),
}
var resp ModifyCloudDocV2Resp
_, err = d.andAlbumRequest(pathname, data, &resp)
if err != nil {
return err
}
if resp.Result.ResultCode != "0" {
return fmt.Errorf("failed to rename family folder: %s", resp.Result.ResultDesc)
}
return nil
} else {
data = base.Json{
"contentID": srcObj.GetID(),
@@ -421,6 +491,33 @@ func (d *Yun139) Copy(ctx context.Context, srcObj, dstDir model.Obj) error {
}
pathname := "/orchestration/personalCloud/batchOprTask/v1.0/createBatchOprTask"
_, err = d.post(pathname, data, nil)
case MetaGroup:
err = d.handleMetaGroupCopy(ctx, srcObj, dstDir)
case MetaFamily:
pathname := "/copyContentCatalog"
var sourceContentIDs []string
var sourceCatalogIDs []string
if srcObj.IsDir() {
sourceCatalogIDs = append(sourceCatalogIDs, srcObj.GetID())
} else {
sourceContentIDs = append(sourceContentIDs, srcObj.GetID())
}
body := base.Json{
"commonAccountInfo": base.Json{
"accountType": "1",
"accountUserId": d.ref.UserDomainID,
},
"destCatalogID": dstDir.GetID(),
"destCloudID": d.CloudID,
"sourceCatalogIDs": sourceCatalogIDs,
"sourceCloudID": d.CloudID,
"sourceContentIDs": sourceContentIDs,
}
var resp base.Json // Assuming a generic JSON response for success/failure
_, err = d.andAlbumRequest(pathname, body, &resp)
// For now, we assume no error means success.
default:
err = errs.NotImplement
}
@@ -680,6 +777,8 @@ func (d *Yun139) Put(ctx context.Context, dstDir model.Obj, stream model.FileStr
return nil
case MetaPersonal:
fallthrough
case MetaGroup:
fallthrough
case MetaFamily:
// 处理冲突
// 获取文件列表
@@ -727,12 +826,17 @@ func (d *Yun139) Put(ctx context.Context, dstDir model.Obj, stream model.FileStr
},
}
pathname := "/orchestration/personalCloud/uploadAndDownload/v1.0/pcUploadFileRequest"
if d.isFamily() {
if d.isFamily() || d.Addition.Type == MetaGroup {
uploadPath := path.Join(dstDir.GetPath(), dstDir.GetID())
// if dstDir is root folder
if dstDir.GetID() == d.RootFolderID {
uploadPath = d.RootPath
}
data = d.newJson(base.Json{
"fileCount": 1,
"manualRename": 2,
"operation": 0,
"path": path.Join(dstDir.GetPath(), dstDir.GetID()),
"path": uploadPath,
"seqNo": random.String(32), // 序列号不能为空
"totalSize": reportSize,
"uploadContentList": []base.Json{{
@@ -744,6 +848,7 @@ func (d *Yun139) Put(ctx context.Context, dstDir model.Obj, stream model.FileStr
pathname = "/orchestration/familyCloud-rebuild/content/v1.0/getFileUploadURL"
}
var resp UploadResp
log.Debugf("[139] upload request body: %+v", data)
_, err = d.post(pathname, data, &resp)
if err != nil {
return err
@@ -839,7 +944,7 @@ func (d *Yun139) GetDetails(ctx context.Context) (*model.StorageDetails, error)
if d.UserDomainID == "" {
return nil, errs.NotImplement
}
var total, free uint64
var total, used uint64
if d.isFamily() {
diskInfo, err := d.getFamilyDiskInfo(ctx)
if err != nil {
@@ -854,7 +959,7 @@ func (d *Yun139) GetDetails(ctx context.Context) (*model.StorageDetails, error)
return nil, fmt.Errorf("failed convert used size into integer: %+v", err)
}
total = totalMb * 1024 * 1024
free = total - (usedMb * 1024 * 1024)
used = usedMb * 1024 * 1024
} else {
diskInfo, err := d.getPersonalDiskInfo(ctx)
if err != nil {
@@ -869,13 +974,10 @@ func (d *Yun139) GetDetails(ctx context.Context) (*model.StorageDetails, error)
return nil, fmt.Errorf("failed convert free size into integer: %+v", err)
}
total = totalMb * 1024 * 1024
free = freeMb * 1024 * 1024
used = total - (freeMb * 1024 * 1024)
}
return &model.StorageDetails{
DiskUsage: model.DiskUsage{
TotalSpace: total,
FreeSpace: free,
},
DiskUsage: driver.DiskUsageFromUsedAndTotal(used, total),
}, nil
}
+3
View File
@@ -8,6 +8,9 @@ import (
type Addition struct {
//Account string `json:"account" required:"true"`
Authorization string `json:"authorization" type:"text" required:"true"`
Username string `json:"username" required:"true"`
Password string `json:"password" required:"true" secret:"true"`
MailCookies string `json:"mail_cookies" required:"true" type:"text" help:"Cookies from mail.139.com used for login authentication."`
driver.RootID
Type string `json:"type" type:"select" options:"personal_new,family,group,personal" default:"personal_new"`
CloudID string `json:"cloud_id"`
+59
View File
@@ -329,3 +329,62 @@ type FamilyDiskInfoResp struct {
DiskSize string `json:"diskSize"`
} `json:"data"`
}
type AndAlbumUploadResp struct {
Result struct {
ResultCode string `json:"resultCode"`
ResultDesc string `json:"resultDesc"`
} `json:"result"`
UploadResult struct {
UploadTaskID string `json:"uploadTaskID"`
RedirectionURL string `json:"redirectionUrl"`
NewContentIDList []struct {
ContentID string `json:"contentID"`
ContentName string `json:"contentName"`
} `json:"newContentIDList"`
} `json:"uploadResult"`
}
type ModifyCloudDocV2Req struct {
CatalogType int `json:"catalogType"`
CloudID string `json:"cloudID"`
CommonAccountInfo struct {
Account string `json:"account"`
AccountType string `json:"accountType"`
} `json:"commonAccountInfo"`
DocLibName string `json:"docLibName"`
DocLibraryID string `json:"docLibraryID"`
Path string `json:"path"`
}
type ModifyCloudDocV2Resp struct {
Result struct {
ResultCode string `json:"resultCode"`
ResultDesc string `json:"resultDesc"`
} `json:"result"`
}
type CreateBatchOprTaskReq struct {
CatalogList []string `json:"catalogList"`
CommonAccountInfo struct {
Account string `json:"account"`
AccountType string `json:"accountType"`
} `json:"commonAccountInfo"`
ContentList []string `json:"contentList"`
DestCatalogID string `json:"destCatalogID"`
DestGroupID string `json:"destGroupID"`
DestPath string `json:"destPath"`
DestType int `json:"destType"`
SourceCatalogType int `json:"sourceCatalogType"`
SourceCloudID string `json:"sourceCloudID"`
SourceType int `json:"sourceType"`
TaskType int `json:"taskType"`
}
type CreateBatchOprTaskResp struct {
Result struct {
ResultCode string `json:"resultCode"`
ResultDesc string `json:"resultDesc"`
} `json:"result"`
TaskID string `json:"taskID"`
}
+702 -8
View File
@@ -1,14 +1,22 @@
package _139
import (
"bytes"
"context"
"crypto/aes"
"crypto/cipher"
"crypto/md5"
crypto_rand "crypto/rand"
"crypto/sha1"
"encoding/base64"
"encoding/hex"
"errors"
"fmt"
"io"
"net/http"
"net/url"
"path"
"regexp"
"sort"
"strconv"
"strings"
@@ -25,6 +33,11 @@ import (
log "github.com/sirupsen/logrus"
)
const (
KEY_HEX_1 = "73634235495062495331515373756c734e7253306c673d3d" // 第一层 AES 解密密钥
KEY_HEX_2 = "7150714477323633586746674c337538" // 第二层 AES 解密密钥
)
// do others that not defined in Driver interface
func (d *Yun139) isFamily() bool {
return d.Type == "family"
@@ -96,12 +109,16 @@ func (d *Yun139) refreshToken() error {
SetBody(reqBody).
SetResult(&resp).
Post(url)
if err != nil {
return err
}
if resp.Return != "0" {
return fmt.Errorf("failed to refresh token: %s", resp.Desc)
if err != nil || resp.Return != "0" {
log.Warnf("139yun: failed to refresh token with old token: %v, desc: %s. trying to login with password.", err, resp.Desc)
newAuth, loginErr := d.loginWithPassword()
log.Debugf("newAuth: Ok: %s", newAuth)
if loginErr != nil {
return fmt.Errorf("failed to login with password after refresh failed: %w", loginErr)
}
return nil
}
d.Authorization = base64.StdEncoding.EncodeToString([]byte(splits[0] + ":" + splits[1] + ":" + resp.Token))
op.MustSaveDriverStorage(d)
return nil
@@ -146,10 +163,29 @@ func (d *Yun139) request(url string, method string, callback base.ReqCallback, r
var e BaseResp
req.SetResult(&e)
log.Debugf("[139] request: %s %s, body: %s", method, url, string(body))
res, err := req.Execute(method, url)
log.Debugln(res.String())
if err != nil {
log.Debugf("[139] request error: %v", err)
return nil, err
}
log.Debugf("[139] response body: %s", res.String())
if !e.Success {
return nil, errors.New(e.Message)
// Always try to unmarshal to the specific response type first if 'resp' is provided.
if resp != nil {
err = utils.Json.Unmarshal(res.Body(), resp)
if err != nil {
log.Debugf("[139] failed to unmarshal response to specific type: %v", err)
return nil, err // Return unmarshal error
}
if createBatchOprTaskResp, ok := resp.(*CreateBatchOprTaskResp); ok {
log.Debugf("[139] CreateBatchOprTaskResp.Result.ResultCode: %s", createBatchOprTaskResp.Result.ResultCode)
if createBatchOprTaskResp.Result.ResultCode == "0" {
goto SUCCESS_PROCESS
}
}
}
return nil, errors.New(e.Message) // Fallback to original error if not handled
}
if resp != nil {
err = utils.Json.Unmarshal(res.Body(), resp)
@@ -157,6 +193,7 @@ func (d *Yun139) request(url string, method string, callback base.ReqCallback, r
return nil, err
}
}
SUCCESS_PROCESS:
return res.Body(), nil
}
@@ -311,6 +348,9 @@ func (d *Yun139) familyGetFiles(catalogID string) ([]model.Obj, error) {
return nil, err
}
path := resp.Data.Path
if catalogID == d.RootFolderID {
d.RootPath = path
}
for _, catalog := range resp.Data.CloudCatalogList {
f := model.Object{
ID: catalog.CatalogID,
@@ -366,6 +406,9 @@ func (d *Yun139) groupGetFiles(catalogID string) ([]model.Obj, error) {
return nil, err
}
path := resp.Data.GetGroupContentResult.ParentCatalogID
if catalogID == d.RootFolderID {
d.RootPath = path
}
for _, catalog := range resp.Data.GetGroupContentResult.CatalogList {
f := model.Object{
ID: catalog.CatalogID,
@@ -494,11 +537,13 @@ func (d *Yun139) personalRequest(pathname string, method string, callback base.R
var e BaseResp
req.SetResult(&e)
log.Debugf("[139] personal request: %s %s, body: %s", method, url, string(body))
res, err := req.Execute(method, url)
if err != nil {
log.Debugf("[139] personal request error: %v", err)
return nil, err
}
log.Debugln(res.String())
log.Debugf("[139] personal response body: %s", res.String())
if !e.Success {
return nil, errors.New(e.Message)
}
@@ -517,6 +562,13 @@ func (d *Yun139) personalPost(pathname string, data interface{}, resp interface{
}, resp)
}
func (d *Yun139) isboPost(pathname string, data interface{}, resp interface{}) ([]byte, error) {
url := "https://group.yun.139.com/hcy/mutual/adapter" + pathname
return d.request(url, http.MethodPost, func(req *resty.Request) {
req.SetBody(data)
}, resp)
}
func getPersonalTime(t string) time.Time {
stamp, err := time.ParseInLocation("2006-01-02T15:04:05.999-07:00", t, utils.CNLoc)
if err != nil {
@@ -703,3 +755,645 @@ func (d *Yun139) getFamilyDiskInfo(ctx context.Context) (*FamilyDiskInfoResp, er
}
return &resp, nil
}
func getMd5(dataStr string) string {
hash := md5.Sum([]byte(dataStr))
return fmt.Sprintf("%x", hash)
}
func (d *Yun139) step1_password_login() (string, error) {
log.Debugf("--- 执行步骤 1: 登录 API ---")
loginURL := "https://mail.10086.cn/Login/Login.ashx"
// 密码 SHA1 哈希
hashedPassword := sha1Hash(fmt.Sprintf("fetion.com.cn:%s", d.Password))
log.Debugf("DEBUG: 原始密码: %s", d.Password)
log.Debugf("DEBUG: SHA1 输入: fetion.com.cn:%s", d.Password)
log.Debugf("DEBUG: 生成的 Password 哈希: %s", hashedPassword)
cguid := strconv.FormatInt(time.Now().UnixMilli(), 10) // 随机生成 cguid
loginHeaders := map[string]string{
"accept": "text/html,application/xhtml+xml,application/xml;q=0.9,image/avif,image/webp,image/apng,*/*;q=0.8,application/signed-exchange;v=b3;q=0.7",
"accept-language": "zh-CN,zh;q=0.9,zh-TW;q=0.8,en-US;q=0.7,en;q=0.6,en-GB;q=0.5",
"cache-control": "max-age=0",
"content-type": "application/x-www-form-urlencoded",
"dnt": "1",
"origin": "https://mail.10086.cn",
"priority": "u=0, i",
"referer": fmt.Sprintf("https://mail.10086.cn/default.html?&s=1&v=0&u=%s&m=1&ec=S001&resource=indexLogin&clientid=1003&auto=on&cguid=%s&mtime=45", base64.StdEncoding.EncodeToString([]byte(d.Username)), cguid),
"sec-ch-ua": "\"Microsoft Edge\";v=\"141\", \"Not?A_Brand\";v=\"8\", \"Chromium\";v=\"141\"",
"sec-ch-ua-mobile": "?0",
"sec-ch-ua-platform": "\"Windows\"",
"sec-fetch-dest": "document",
"sec-fetch-mode": "navigate",
"sec-fetch-site": "same-origin",
"sec-fetch-user": "?1",
"upgrade-insecure-requests": "1",
"user-agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/141.0.0.0 Safari/537.36 Edg/141.0.0.0",
"Cookie": d.MailCookies,
}
loginData := url.Values{}
loginData.Set("UserName", d.Username)
loginData.Set("passOld", "")
loginData.Set("auto", "on")
loginData.Set("Password", hashedPassword)
loginData.Set("webIndexPagePwdLogin", "1")
loginData.Set("pwdType", "1")
loginData.Set("clientId", "1003")
loginData.Set("authType", "2")
log.Debugf("DEBUG: 登录请求 URL: %s", loginURL)
log.Debugf("DEBUG: 登录请求 Headers: %+v", loginHeaders)
log.Debugf("DEBUG: 登录请求 Body: %s", loginData.Encode())
// 设置客户端不跟随重定向
client := base.RestyClient.SetRedirectPolicy(resty.NoRedirectPolicy())
res, err := client.R().
SetHeaders(loginHeaders).
SetFormDataFromValues(loginData).
Post(loginURL)
if err != nil {
// 如果是重定向错误,则不作为失败处理,因为我们禁止了自动重定向
if res != nil && res.StatusCode() >= 300 && res.StatusCode() < 400 {
log.Debugf("DEBUG: 登录响应 Status Code: %d (Redirect)", res.StatusCode())
} else {
return "", fmt.Errorf("step1 login request failed: %w", err)
}
} else {
log.Debugf("DEBUG: 登录响应 Status Code: %d", res.StatusCode())
}
// 恢复客户端的默认重定向策略,以免影响后续请求
base.RestyClient.SetRedirectPolicy(resty.FlexibleRedirectPolicy(10))
log.Debugf("DEBUG: 登录响应 Headers: %+v", res.Header())
var sid, extractedCguid string
// 从 Location 头部提取 sid 和 cguid
locationHeader := res.Header().Get("Location")
if locationHeader != "" {
sidMatch := regexp.MustCompile(`sid=([^&]+)`).FindStringSubmatch(locationHeader)
cguidMatch := regexp.MustCompile(`cguid=([^&]+)`).FindStringSubmatch(locationHeader)
if len(sidMatch) > 1 {
sid = sidMatch[1]
log.Debugf("DEBUG: 从 Location 提取到 sid: %s", sid)
}
if len(cguidMatch) > 1 {
extractedCguid = cguidMatch[1]
log.Debugf("DEBUG: 从 Location 提取到 cguid: %s", extractedCguid)
}
}
// 如果 Location 中没有,尝试从 Set-Cookie 中提取
if sid == "" || extractedCguid == "" {
setCookieHeaders := res.Header().Values("Set-Cookie")
for _, cookieStr := range setCookieHeaders {
ssoSidMatch := regexp.MustCompile(`Os_SSo_Sid=([^;]+)`).FindStringSubmatch(cookieStr)
cookieCguidMatch := regexp.MustCompile(`cguid=([^;]+)`).FindStringSubmatch(cookieStr)
if len(ssoSidMatch) > 1 && sid == "" {
sid = ssoSidMatch[1]
log.Debugf("DEBUG: 从 Set-Cookie 提取到 sid: %s", sid)
}
if len(cookieCguidMatch) > 1 && extractedCguid == "" {
extractedCguid = cookieCguidMatch[1]
log.Debugf("DEBUG: 从 Set-Cookie 提取到 cguid: %s", extractedCguid)
}
}
}
if sid == "" || extractedCguid == "" {
return "", errors.New("failed to extract sid or cguid from login response")
}
// 提取并记录 cookies
loginUrlObj, _ := url.Parse(loginURL)
cookies := base.RestyClient.GetClient().Jar.Cookies(loginUrlObj)
var cookieStrings []string
for _, cookie := range cookies {
cookieStrings = append(cookieStrings, cookie.Name+"="+cookie.Value)
}
cookieStr := strings.Join(cookieStrings, "; ")
log.Debugf("DEBUG: 提取到的 Cookies: %s", cookieStr)
d.MailCookies = cookieStr
return sid, nil
}
func (d *Yun139) step2_get_single_token(sid string) (string, error) {
log.Debugf("\n--- 执行步骤 2: 换artifact API ---")
cguid := strconv.FormatInt(time.Now().UnixMilli(), 10)
exchangeArtifactURL := fmt.Sprintf("https://smsrebuild1.mail.10086.cn/setting/s?func=%s&sid=%s&cguid=%s", url.QueryEscape("umc:getArtifact"), sid, cguid)
// 从 MailCookies 中提取 RMKEY
var rmkey string
cookies := strings.Split(d.MailCookies, ";")
for _, cookie := range cookies {
cookie = strings.TrimSpace(cookie)
if strings.HasPrefix(cookie, "RMKEY=") {
rmkey = cookie
break
}
}
if rmkey == "" {
return "", errors.New("RMKEY not found in MailCookies")
}
exchangePassidHeaders := map[string]string{
"Host": "smsrebuild1.mail.10086.cn",
"Cookie": rmkey,
"Content-Type": "text/xml; charset=utf-8",
"Accept-Encoding": "gzip",
"User-Agent": "okhttp/4.12.0",
}
log.Debugf("DEBUG: 换passid 请求 URL: %s", exchangeArtifactURL)
log.Debugf("DEBUG: 换passid 请求 Headers: %+v", exchangePassidHeaders)
res, err := base.RestyClient.R().
SetHeaders(exchangePassidHeaders).
Post(exchangeArtifactURL)
if err != nil {
return "", fmt.Errorf("step2 exchange artifact request failed: %w", err)
}
log.Debugf("DEBUG: 换passid 响应 Status Code: %d", res.StatusCode())
log.Debugf("DEBUG: 换passid 响应 Headers: %+v", res.Header())
log.Debugf("DEBUG: 换passid 响应 Body: %s...", res.String()[:min(len(res.String()), 500)])
dycpwd := jsoniter.Get(res.Body(), "var", "artifact").ToString()
if dycpwd == "" {
return "", errors.New("failed to extract dycpwd from artifact exchange response")
}
log.Debugf("DEBUG: 提取到 dycpwd: %s", dycpwd)
return dycpwd, nil
}
// --- 辅助函数:加密/解密 ---
// sha1Hash 计算 SHA1 哈希值,返回十六进制字符串。
func sha1Hash(data string) string {
h := sha1.New()
h.Write([]byte(data))
return hex.EncodeToString(h.Sum(nil))
}
// pkcs7_pad PKCS7 填充
func pkcs7_pad(data []byte, blockSize int) []byte {
padding := blockSize - len(data)%blockSize
padtext := bytes.Repeat([]byte{byte(padding)}, padding)
return append(data, padtext...)
}
// pkcs7_unpad PKCS7 去填充
func pkcs7_unpad(data []byte) ([]byte, error) {
length := len(data)
if length == 0 {
return nil, errors.New("pkcs7: data is empty")
}
unpadding := int(data[length-1])
if unpadding > length {
return nil, errors.New("pkcs7: invalid padding")
}
return data[:(length - unpadding)], nil
}
// aes_ecb_decrypt AES/ECB/Pkcs7 解密,输入为十六进制字符串。
func aes_ecb_decrypt(ciphertext []byte, key []byte) ([]byte, error) {
block, err := aes.NewCipher(key)
if err != nil {
return nil, err
}
if len(ciphertext)%block.BlockSize() != 0 {
return nil, errors.New("AES ECB decrypt: ciphertext is not a multiple of the block size")
}
decrypted := make([]byte, len(ciphertext))
blockSize := block.BlockSize()
for bs, be := 0, blockSize; bs < len(ciphertext); bs, be = bs+blockSize, be+blockSize {
block.Decrypt(decrypted[bs:be], ciphertext[bs:be])
}
return pkcs7_unpad(decrypted)
}
// 以下提供 camelCase 的 AES CBC 加解密,供文件中其它位置调用(并支持传入 IV)。
func aesCbcEncrypt(plaintext []byte, key []byte, iv []byte) ([]byte, error) {
block, err := aes.NewCipher(key)
if err != nil {
return nil, err
}
if len(iv) != block.BlockSize() {
return nil, fmt.Errorf("aesCbcEncrypt: iv length %d does not match block size %d", len(iv), block.BlockSize())
}
padded := pkcs7_pad(plaintext, block.BlockSize())
ciphertext := make([]byte, len(padded))
mode := cipher.NewCBCEncrypter(block, iv)
mode.CryptBlocks(ciphertext, padded)
return ciphertext, nil
}
func aesCbcDecrypt(ciphertext []byte, key []byte, iv []byte) ([]byte, error) {
block, err := aes.NewCipher(key)
if err != nil {
return nil, err
}
if len(iv) != block.BlockSize() {
return nil, fmt.Errorf("aesCbcDecrypt: iv length %d does not match block size %d", len(iv), block.BlockSize())
}
if len(ciphertext)%block.BlockSize() != 0 {
return nil, errors.New("aesCbcDecrypt: ciphertext is not a multiple of the block size")
}
decrypted := make([]byte, len(ciphertext))
mode := cipher.NewCBCDecrypter(block, iv)
mode.CryptBlocks(decrypted, ciphertext)
return pkcs7_unpad(decrypted)
}
// sortedJsonStringify 对 JSON 对象进行排序并字符串化。
func sortedJsonStringify(obj interface{}) (string, error) {
if obj == nil {
return "null", nil
}
switch v := obj.(type) {
case string:
// 尝试解析为 JSON,如果成功则递归处理
var parsed interface{}
if err := jsoniter.Unmarshal([]byte(v), &parsed); err == nil {
return sortedJsonStringify(parsed)
}
// 如果不是 JSON 字符串,则直接返回 JSON 字符串化的结果
return jsoniter.MarshalToString(v)
case int, float64, bool:
return fmt.Sprintf("%v", v), nil
case []interface{}:
var items []string
for _, item := range v {
s, err := sortedJsonStringify(item)
if err != nil {
return "", err
}
items = append(items, s)
}
return fmt.Sprintf("[%s]", strings.Join(items, ",")), nil
case map[string]interface{}:
sortedKeys := make([]string, 0, len(v))
for key := range v {
sortedKeys = append(sortedKeys, key)
}
sort.Strings(sortedKeys)
var pairs []string
for _, key := range sortedKeys {
value := v[key]
s, err := sortedJsonStringify(value)
if err != nil {
return "", err
}
// Use jsoniter.MarshalToString for the key to ensure it's quoted correctly
keyStr, err := jsoniter.MarshalToString(key)
if err != nil {
return "", err
}
pairs = append(pairs, fmt.Sprintf("%s:%s", keyStr, s))
}
return fmt.Sprintf("{%s}", strings.Join(pairs, ",")), nil
default:
// Fallback for other types, e.g., numbers, booleans, or unhandled complex types
// Use jsoniter's default marshalling for these
return jsoniter.MarshalToString(v)
}
}
// yun139EncryptedRequest handles the common encrypted request/response flow.
func (d *Yun139) yun139EncryptedRequest(url string, body interface{}, headers map[string]string, aesKeyHex string, resp interface{}) ([]byte, error) {
// 1. Decode AES key
aesKey, err := hex.DecodeString(aesKeyHex)
if err != nil {
return nil, fmt.Errorf("yun139EncryptedRequest: failed to decode AES key: %w", err)
}
// 2. Marshal and sort the request body
sortedJson, err := sortedJsonStringify(body)
if err != nil {
return nil, fmt.Errorf("yun139EncryptedRequest: failed to marshal and sort body: %w", err)
}
log.Debugf("yun139EncryptedRequest: Request Body (plaintext): %s", sortedJson)
// 3. Encrypt the body using AES/CBC
iv := make([]byte, 16) // 16 bytes for AES-128
if _, err := crypto_rand.Read(iv); err != nil {
return nil, fmt.Errorf("yun139EncryptedRequest: failed to generate IV: %w", err)
}
encryptedBody, err := aesCbcEncrypt([]byte(sortedJson), aesKey, iv)
if err != nil {
return nil, fmt.Errorf("yun139EncryptedRequest: failed to encrypt body: %w", err)
}
payload := base64.StdEncoding.EncodeToString(append(iv, encryptedBody...))
// 4. Make the request
res, err := base.RestyClient.R().
SetHeaders(headers).
SetBody(payload).
Post(url)
if err != nil {
return nil, fmt.Errorf("yun139EncryptedRequest: http request failed: %w", err)
}
if res.StatusCode() != 200 {
return nil, fmt.Errorf("yun139EncryptedRequest: unexpected status code %d: %s", res.StatusCode(), res.String())
}
// 5. Decrypt the response
respBody := res.Body()
var decryptedBytes []byte
if len(respBody) > 0 && respBody[0] == '{' {
log.Warnf("yun139EncryptedRequest: received a plain JSON response, not an encrypted string. Body: %s", string(respBody))
decryptedBytes = respBody
} else {
decodedResp, err := base64.StdEncoding.DecodeString(string(respBody))
if err != nil {
return nil, fmt.Errorf("yun139EncryptedRequest: response base64 decode failed: %w. Body: '%s'", err, string(respBody))
}
if len(decodedResp) < 16 {
return nil, fmt.Errorf("yun139EncryptedRequest: decoded response is too short to be encrypted. Length: %d", len(decodedResp))
}
respIv := decodedResp[:16]
respCiphertext := decodedResp[16:]
decryptedBytes, err = aesCbcDecrypt(respCiphertext, aesKey, respIv)
if err != nil {
return nil, fmt.Errorf("yun139EncryptedRequest: response aes decrypt failed: %w", err)
}
}
log.Debugf("yun139EncryptedRequest: Response Body (decrypted): %s", string(decryptedBytes))
// 6. Unmarshal to the final response struct
if resp != nil {
err = utils.Json.Unmarshal(decryptedBytes, resp)
if err != nil {
return nil, fmt.Errorf("yun139EncryptedRequest: failed to unmarshal decrypted response: %w", err)
}
}
return decryptedBytes, nil
}
func (d *Yun139) step3_third_party_login(dycpwd string) (string, error) {
log.Debugf("\n--- 执行步骤 3: 单点登录 API ---")
ssoLoginURL := "https://user-njs.yun.139.com/user/thirdlogin"
// 构建原始请求体
ssoRequestBodyRaw := base.Json{
"clientkey_decrypt": "l3TryM&Q+X7@dzwk)qP",
"clienttype": "886",
"cpid": "507",
"dycpwd": dycpwd,
"extInfo": base.Json{"ifOpenAccount": "0"},
"loginMode": "0",
"msisdn": d.Username,
"pintype": "13",
"secinfo": strings.ToUpper(sha1Hash(fmt.Sprintf("fetion.com.cn:%s", dycpwd))),
"version": "20250901",
}
ssoLoginHeaders := map[string]string{
"hcy-cool-flag": "1",
"x-huawei-channelSrc": "10246600",
"x-sdk-channelSrc": "",
"x-MM-Source": "0",
"x-UserAgent": "android|23116PN5BC|android15|1.2.6|||1440x3200|10246600",
"x-DeviceInfo": "4|127.0.0.1|5|1.2.6|Xiaomi|23116PN5BC||02-00-00-00-00-00|android 15|1440x3200|android|||",
"Content-Type": "text/plain;charset=UTF-8",
"Host": "user-njs.yun.139.com",
"Connection": "Keep-Alive",
"Accept-Encoding": "gzip",
"User-Agent": "okhttp/3.12.2",
}
// 使用通用加密请求函数
decryptedLayer1StrBytes, err := d.yun139EncryptedRequest(ssoLoginURL, ssoRequestBodyRaw, ssoLoginHeaders, KEY_HEX_1, nil)
if err != nil {
return "", fmt.Errorf("step3 encrypted request failed: %w", err)
}
hexInner := jsoniter.Get(decryptedLayer1StrBytes, "data").ToString()
if hexInner == "" {
return "", errors.New("missing data field in first layer decryption result")
}
log.Debugf("DEBUG: 第一层解密提取到 hex_inner: %s...", hexInner[:min(len(hexInner), 50)])
// 第二层解密
key2, err := hex.DecodeString(KEY_HEX_2)
if err != nil {
return "", fmt.Errorf("failed to decode KEY_HEX_2: %w", err)
}
hexInnerBytes, err := hex.DecodeString(hexInner)
if err != nil {
return "", fmt.Errorf("failed to decode hex_inner: %w", err)
}
finalJsonStrBytes, err := aes_ecb_decrypt(hexInnerBytes, key2)
if err != nil {
return "", fmt.Errorf("step3 response layer2 aes ecb decrypt failed: %w", err)
}
log.Debugf("DEBUG: 最终解密结果: %s", string(finalJsonStrBytes))
// 提取 authToken
authToken := jsoniter.Get(finalJsonStrBytes, "authToken").ToString()
if authToken == "" {
return "", errors.New("failed to extract authToken from final decryption result")
}
log.Debugf("DEBUG: 提取到 authToken: %s", authToken)
// 提取 account 和 userDomainId
account := jsoniter.Get(finalJsonStrBytes, "account").ToString()
userDomainId := jsoniter.Get(finalJsonStrBytes, "userDomainId").ToString()
if account == "" || userDomainId == "" {
return "", errors.New("failed to extract account or userDomainId from final decryption result")
}
d.UserDomainID = userDomainId
newAuthorization := base64.StdEncoding.EncodeToString([]byte(fmt.Sprintf("pc:%s:%s", account, authToken)))
return newAuthorization, nil
}
func (d *Yun139) loginWithPassword() (string, error) {
if d.Username == "" || d.Password == "" || d.MailCookies == "" {
return "", errors.New("username, password or mail_cookies is empty")
}
passId, err := d.step1_password_login()
if err != nil {
return "", err
}
log.Infof("Step 1 success, passId: %s", passId)
token, err := d.step2_get_single_token(passId)
if err != nil {
return "", err
}
log.Infof("Step 2 success, token: %s", token)
newAuth, err := d.step3_third_party_login(token)
if err != nil {
return "", err
}
log.Infof("Step 3 success, new authorization generated.")
d.Authorization = newAuth // Ensure Authorization is also updated before saving
op.MustSaveDriverStorage(d)
return newAuth, nil
}
func (d *Yun139) andAlbumRequest(pathname string, body interface{}, resp interface{}) ([]byte, error) {
url := "https://group.yun.139.com/hcy/family/adapter/andAlbum/openApi" + pathname
headers := map[string]string{
"Host": "group.yun.139.com",
"authorization": "Basic " + d.getAuthorization(),
"x-svctype": "2",
"hcy-cool-flag": "1",
"api-version": "v2",
"x-huawei-channelsrc": "10246600",
"x-sdk-channelsrc": "",
"x-mm-source": "0",
"x-deviceinfo": "1|127.0.0.1|1|12.3.2|Xiaomi|23116PN5BC||02-00-00-00-00-00|android 15|1440x3200|android|zh||||032|0|", //重要参数
"content-type": "application/json; charset=utf-8",
"user-agent": "okhttp/4.11.0",
"accept-encoding": "gzip",
}
return d.yun139EncryptedRequest(url, body, headers, KEY_HEX_1, resp)
}
func (d *Yun139) handleMetaGroupCopy(ctx context.Context, srcObj, dstDir model.Obj) error {
pathname := "/copyContentCatalog"
var sourceContentIDs []string
var sourceCatalogIDs []string
if srcObj.IsDir() {
sourceCatalogIDs = append(sourceCatalogIDs, path.Join("root:/", srcObj.GetPath(), srcObj.GetID()))
} else {
sourceContentIDs = append(sourceContentIDs, path.Join("root:/", srcObj.GetPath(), srcObj.GetID()))
}
destCatalogID := path.Join("root:/", dstDir.GetPath(), dstDir.GetID())
log.Debugf("[139Yun Group Copy] srcObj ID: %s, srcObj Path: %s, dstDir ID: %s, dstDir Path: %s, destCatalogID: %s", srcObj.GetID(), srcObj.GetPath(), dstDir.GetID(), dstDir.GetPath(), destCatalogID)
body := base.Json{
"commonAccountInfo": base.Json{
"accountType": "1",
"accountUserId": d.UserDomainID,
},
"destCatalogID": destCatalogID,
"destCloudID": d.CloudID,
"sourceCatalogIDs": sourceCatalogIDs,
"sourceCloudID": d.CloudID,
"sourceContentIDs": sourceContentIDs,
}
var resp base.Json
_, err := d.andAlbumRequest(pathname, body, &resp)
return err
}
// getGroupRootByCloudID 查询 group 上层信息,优先返回 parentCatalogID,回退到 catalogList[0].path
func (d *Yun139) getGroupRootByCloudID(cloudID string) (string, error) {
pathname := "/orchestration/group-rebuild/catalog/v1.0/queryGroupContentList"
body := base.Json{
"groupID": cloudID,
"commonAccountInfo": base.Json{
"account": d.getAccount(),
"accountType": 1,
},
"pageInfo": base.Json{
"pageNum": 1,
"pageSize": 1,
},
}
var resp base.Json
_, err := d.post(pathname, body, &resp)
if err != nil {
return "", err
}
dataObj, _ := resp["data"].(map[string]interface{})
if dataObj == nil {
return "", fmt.Errorf("invalid group response data")
}
if gcr, ok := dataObj["getGroupContentResult"].(map[string]interface{}); ok {
if pid, ok := gcr["parentCatalogID"].(string); ok && pid != "" {
return pid, nil
}
if cl, ok := gcr["catalogList"].([]interface{}); ok && len(cl) > 0 {
if first, ok := cl[0].(map[string]interface{}); ok {
if p, ok := first["path"].(string); ok && p != "" {
return p, nil
}
}
}
}
return "", fmt.Errorf("no root found in group response")
}
// getFamilyRootPath 查询 family 的上层 path(data.path)
// 返回值已去除前缀 "root:/"(或 "root:"),直接返回纯 ID 或 path 部分,便于持久化为 RootFolderID。
func (d *Yun139) getFamilyRootPath(cloudID string) (string, error) {
// 使用 v1.2 接口(代码日志中已有该请求),pageSize 取 1 足够获取 path 字段
pathname := "/orchestration/familyCloud-rebuild/content/v1.2/queryContentList"
body := base.Json{
"catalogID": "",
"catalogType": 3,
"cloudID": cloudID,
"cloudType": 1,
"commonAccountInfo": base.Json{
"account": d.getAccount(),
"accountType": 1,
},
"contentSortType": 0,
"pageInfo": base.Json{
"pageNum": 1,
"pageSize": 1,
},
"sortDirection": 1,
}
var resp base.Json
_, err := d.post(pathname, body, &resp)
if err != nil {
return "", err
}
dataObj, _ := resp["data"].(map[string]interface{})
if dataObj == nil {
return "", fmt.Errorf("invalid family response data")
}
// helper to strip "root:/" or "root:" prefix
stripRoot := func(s string) string {
s = strings.TrimSpace(s)
s = strings.TrimPrefix(s, "root:/")
s = strings.TrimPrefix(s, "root:")
return s
}
if p, ok := dataObj["path"].(string); ok && p != "" {
return stripRoot(p), nil
}
// 回退:有时 path 在 cloudCatalogList.catalogList 中
if cl, ok := dataObj["cloudCatalogList"].([]interface{}); ok && len(cl) > 0 {
if first, ok := cl[0].(map[string]interface{}); ok {
if p, ok := first["path"].(string); ok && p != "" {
return stripRoot(p), nil
}
}
}
return "", fmt.Errorf("no path found in family response")
}
+1 -4
View File
@@ -200,10 +200,7 @@ func (d *Cloud189) GetDetails(ctx context.Context) (*model.StorageDetails, error
return nil, err
}
return &model.StorageDetails{
DiskUsage: model.DiskUsage{
TotalSpace: capacityInfo.CloudCapacityInfo.TotalSize,
FreeSpace: capacityInfo.CloudCapacityInfo.FreeSize,
},
DiskUsage: driver.DiskUsageFromUsedAndTotal(capacityInfo.CloudCapacityInfo.UsedSize, capacityInfo.CloudCapacityInfo.TotalSize),
}, nil
}
+2 -2
View File
@@ -72,13 +72,13 @@ type CapacityResp struct {
ResMessage string `json:"res_message"`
Account string `json:"account"`
CloudCapacityInfo struct {
FreeSize uint64 `json:"freeSize"`
FreeSize int64 `json:"freeSize"`
MailUsedSize uint64 `json:"mail189UsedSize"`
TotalSize uint64 `json:"totalSize"`
UsedSize uint64 `json:"usedSize"`
} `json:"cloudCapacityInfo"`
FamilyCapacityInfo struct {
FreeSize uint64 `json:"freeSize"`
FreeSize int64 `json:"freeSize"`
TotalSize uint64 `json:"totalSize"`
UsedSize uint64 `json:"usedSize"`
} `json:"familyCapacityInfo"`
+1 -1
View File
@@ -107,7 +107,7 @@ import (
// res, err = d.client.R().
// SetHeaders(map[string]string{
// "lt": lt,
// "User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/87.0.4280.88 Safari/537.36",
// "User-Agent": base.UserAgentNT,
// "Referer": "https://open.e.189.cn/",
// "accept": "application/json;charset=UTF-8",
// }).SetFormData(map[string]string{
+4 -7
View File
@@ -284,18 +284,15 @@ func (y *Cloud189TV) GetDetails(ctx context.Context) (*model.StorageDetails, err
if err != nil {
return nil, err
}
var total, free uint64
var total, used uint64
if y.isFamily() {
total = capacityInfo.FamilyCapacityInfo.TotalSize
free = capacityInfo.FamilyCapacityInfo.FreeSize
used = capacityInfo.FamilyCapacityInfo.UsedSize
} else {
total = capacityInfo.CloudCapacityInfo.TotalSize
free = capacityInfo.CloudCapacityInfo.FreeSize
used = capacityInfo.CloudCapacityInfo.UsedSize
}
return &model.StorageDetails{
DiskUsage: model.DiskUsage{
TotalSpace: total,
FreeSpace: free,
},
DiskUsage: driver.DiskUsageFromUsedAndTotal(used, total),
}, nil
}
+2 -2
View File
@@ -322,13 +322,13 @@ type CapacityResp struct {
ResMessage string `json:"res_message"`
Account string `json:"account"`
CloudCapacityInfo struct {
FreeSize uint64 `json:"freeSize"`
FreeSize int64 `json:"freeSize"`
MailUsedSize uint64 `json:"mail189UsedSize"`
TotalSize uint64 `json:"totalSize"`
UsedSize uint64 `json:"usedSize"`
} `json:"cloudCapacityInfo"`
FamilyCapacityInfo struct {
FreeSize uint64 `json:"freeSize"`
FreeSize int64 `json:"freeSize"`
TotalSize uint64 `json:"totalSize"`
UsedSize uint64 `json:"usedSize"`
} `json:"familyCapacityInfo"`
+4 -7
View File
@@ -416,18 +416,15 @@ func (y *Cloud189PC) GetDetails(ctx context.Context) (*model.StorageDetails, err
if err != nil {
return nil, err
}
var total, free uint64
var total, used uint64
if y.isFamily() {
total = capacityInfo.FamilyCapacityInfo.TotalSize
free = capacityInfo.FamilyCapacityInfo.FreeSize
used = capacityInfo.FamilyCapacityInfo.UsedSize
} else {
total = capacityInfo.CloudCapacityInfo.TotalSize
free = capacityInfo.CloudCapacityInfo.FreeSize
used = capacityInfo.CloudCapacityInfo.UsedSize
}
return &model.StorageDetails{
DiskUsage: model.DiskUsage{
TotalSpace: total,
FreeSpace: free,
},
DiskUsage: driver.DiskUsageFromUsedAndTotal(used, total),
}, nil
}
+2 -2
View File
@@ -415,13 +415,13 @@ type CapacityResp struct {
ResMessage string `json:"res_message"`
Account string `json:"account"`
CloudCapacityInfo struct {
FreeSize uint64 `json:"freeSize"`
FreeSize int64 `json:"freeSize"`
MailUsedSize uint64 `json:"mail189UsedSize"`
TotalSize uint64 `json:"totalSize"`
UsedSize uint64 `json:"usedSize"`
} `json:"cloudCapacityInfo"`
FamilyCapacityInfo struct {
FreeSize uint64 `json:"freeSize"`
FreeSize int64 `json:"freeSize"`
TotalSize uint64 `json:"totalSize"`
UsedSize uint64 `json:"usedSize"`
} `json:"familyCapacityInfo"`
+19 -23
View File
@@ -756,30 +756,24 @@ func (y *Cloud189PC) StreamUpload(ctx context.Context, dstDir model.Obj, file mo
}
partInfo := ""
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, partSize)
if err != nil {
return err
}
silceMd5.Reset()
w, err := utils.CopyWithBuffer(writers, 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))
rateLimitedRd = driver.NewLimitedUploadStream(ctx, reader)
Before: func(ctx context.Context) (err error) {
reader, err = ss.GetSectionReader(offset, partSize)
if err != nil {
return err
}
silceMd5.Reset()
w, err := utils.CopyWithBuffer(writers, 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))
return nil
},
Do: func(ctx context.Context) error {
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 {
@@ -788,11 +782,11 @@ func (y *Cloud189PC) StreamUpload(ctx context.Context, dstDir model.Obj, file mo
// step.4 上传切片
uploadUrl := uploadUrls[0]
_, err = y.put(ctx, uploadUrl.RequestURL, uploadUrl.Headers, false, rateLimitedRd, isFamily)
_, err = y.put(ctx, uploadUrl.RequestURL, uploadUrl.Headers, false, driver.NewLimitedUploadStream(ctx, reader), isFamily)
if err != nil {
return err
}
up(float64(threadG.Success()) * 100 / float64(count))
up(float64(threadG.Success()+1) * 100 / float64(count+1))
return nil
},
After: func(err error) {
@@ -804,6 +798,7 @@ func (y *Cloud189PC) StreamUpload(ctx context.Context, dstDir model.Obj, file mo
if err = threadG.Wait(); err != nil {
return nil, err
}
defer up(100)
if fileMd5 != nil {
fileMd5Hex = strings.ToUpper(hex.EncodeToString(fileMd5.Sum(nil)))
@@ -995,7 +990,7 @@ func (y *Cloud189PC) FastUpload(ctx context.Context, dstDir model.Obj, file mode
return err
}
up(float64(threadG.Success()) * 100 / float64(len(uploadUrls)))
up(float64(threadG.Success()+1) * 100 / float64(len(uploadUrls)+1))
uploadProgress.UploadParts[i] = ""
return nil
})
@@ -1007,6 +1002,7 @@ func (y *Cloud189PC) FastUpload(ctx context.Context, dstDir model.Obj, file mode
}
return nil, err
}
defer up(100)
}
// step.5 提交
+322 -331
View File
@@ -5,6 +5,7 @@ import (
"errors"
"fmt"
"io"
"math/rand"
"net/url"
stdpath "path"
"strings"
@@ -16,6 +17,7 @@ import (
"github.com/OpenListTeam/OpenList/v4/internal/op"
"github.com/OpenListTeam/OpenList/v4/internal/sign"
"github.com/OpenListTeam/OpenList/v4/internal/stream"
"github.com/OpenListTeam/OpenList/v4/pkg/http_range"
"github.com/OpenListTeam/OpenList/v4/pkg/utils"
"github.com/OpenListTeam/OpenList/v4/server/common"
)
@@ -23,10 +25,9 @@ import (
type Alias struct {
model.Storage
Addition
rootOrder []string
pathMap map[string][]string
autoFlatten bool
oneKey string
rootOrder []string
pathMap map[string][]string
root model.Obj
}
func (d *Alias) Config() driver.Config {
@@ -38,9 +39,6 @@ func (d *Alias) GetAddition() driver.Additional {
}
func (d *Alias) Init(ctx context.Context) error {
if d.Paths == "" {
return errors.New("paths is required")
}
paths := strings.Split(d.Paths, "\n")
d.rootOrder = make([]string, 0, len(paths))
d.pathMap = make(map[string][]string)
@@ -50,19 +48,50 @@ func (d *Alias) Init(ctx context.Context) error {
continue
}
k, v := getPair(path)
if _, ok := d.pathMap[k]; !ok {
temp, ok := d.pathMap[k]
if !ok {
d.rootOrder = append(d.rootOrder, k)
}
d.pathMap[k] = append(d.pathMap[k], v)
d.pathMap[k] = append(temp, v)
}
if len(d.pathMap) == 1 {
for k := range d.pathMap {
d.oneKey = k
switch len(d.rootOrder) {
case 0:
return errors.New("paths is required")
case 1:
paths := d.pathMap[d.rootOrder[0]]
roots := make(BalancedObjs, 0, len(paths))
roots = append(roots, &model.Object{
Name: "root",
Path: paths[0],
IsFolder: true,
Modified: d.Modified,
Mask: model.Locked,
})
for _, path := range paths[1:] {
roots = append(roots, &model.Object{
Path: path,
})
}
d.autoFlatten = true
} else {
d.oneKey = ""
d.autoFlatten = false
d.root = roots
default:
d.root = &model.Object{
Name: "root",
Path: "/",
IsFolder: true,
Modified: d.Modified,
Mask: model.ReadOnly,
}
}
if !utils.SliceContains(ValidReadConflictPolicy, d.ReadConflictPolicy) {
d.ReadConflictPolicy = FirstRWP
}
if !utils.SliceContains(ValidWriteConflictPolicy, d.WriteConflictPolicy) {
d.WriteConflictPolicy = DisabledWP
}
if !utils.SliceContains(ValidPutConflictPolicy, d.PutConflictPolicy) {
d.PutConflictPolicy = DisabledWP
}
return nil
}
@@ -70,310 +99,292 @@ func (d *Alias) Init(ctx context.Context) error {
func (d *Alias) Drop(ctx context.Context) error {
d.rootOrder = nil
d.pathMap = nil
d.root = nil
return nil
}
func (d *Alias) Get(ctx context.Context, path string) (model.Obj, error) {
if utils.PathEqual(path, "/") {
return &model.Object{
Name: "Root",
IsFolder: true,
Path: "/",
}, nil
func (d *Alias) GetRoot(ctx context.Context) (model.Obj, error) {
if d.root == nil {
return nil, errs.StorageNotInit
}
root, sub := d.getRootAndPath(path)
dsts, ok := d.pathMap[root]
if !ok {
return d.root, nil
}
// 通过op.Get调用的话,path一定是子路径(/开头)
func (d *Alias) Get(ctx context.Context, path string) (model.Obj, error) {
roots, sub := d.getRootsAndPath(path)
if len(roots) == 0 {
return nil, errs.ObjectNotFound
}
var ret *model.Object
provider := ""
for _, dst := range dsts {
rawPath := stdpath.Join(dst, sub)
for idx, root := range roots {
rawPath := stdpath.Join(root, sub)
obj, err := fs.Get(ctx, rawPath, &fs.GetArgs{NoLog: true})
if err != nil {
continue
}
storage, err := fs.GetStorage(rawPath, &fs.GetStoragesArgs{})
if ret == nil {
ret = &model.Object{
Path: path,
Name: obj.GetName(),
Size: obj.GetSize(),
Modified: obj.ModTime(),
IsFolder: obj.IsDir(),
HashInfo: obj.GetHash(),
}
if !d.ProviderPassThrough || err != nil {
break
}
provider = storage.Config().Name
} else if err != nil || provider != storage.GetStorage().Driver {
provider = ""
break
mask := model.GetObjMask(obj) &^ model.Temp
if sub == "" {
// 根目录
mask |= model.Locked | model.Virtual
}
ret := model.Object{
Path: rawPath,
Name: obj.GetName(),
Size: obj.GetSize(),
Modified: obj.ModTime(),
IsFolder: obj.IsDir(),
HashInfo: obj.GetHash(),
Mask: mask,
}
obj = &ret
if d.ProviderPassThrough && !obj.IsDir() {
if storage, err := fs.GetStorage(rawPath, &fs.GetStoragesArgs{}); err == nil {
obj = &model.ObjectProvider{
Object: ret,
Provider: model.Provider{
Provider: storage.Config().Name,
},
}
}
}
roots = roots[idx+1:]
var objs BalancedObjs
if idx > 0 {
objs = make(BalancedObjs, 0, len(roots)+2)
} else {
objs = make(BalancedObjs, 0, len(roots)+1)
}
objs = append(objs, obj)
if idx > 0 {
objs = append(objs, nil)
}
for _, d := range roots {
objs = append(objs, &tempObj{model.Object{
Path: stdpath.Join(d, sub),
}})
}
return objs, nil
}
if ret == nil {
return nil, errs.ObjectNotFound
}
if provider != "" {
return &model.ObjectProvider{
Object: *ret,
Provider: model.Provider{
Provider: provider,
},
}, nil
}
return ret, nil
return nil, errs.ObjectNotFound
}
func (d *Alias) List(ctx context.Context, dir model.Obj, args model.ListArgs) ([]model.Obj, error) {
path := dir.GetPath()
if utils.PathEqual(path, "/") && !d.autoFlatten {
dirs, ok := dir.(BalancedObjs)
if !ok {
return d.listRoot(ctx, args.WithStorageDetails && d.DetailsPassThrough, args.Refresh), nil
}
root, sub := d.getRootAndPath(path)
dsts, ok := d.pathMap[root]
if !ok {
return nil, errs.ObjectNotFound
}
var objs []model.Obj
for _, dst := range dsts {
tmp, err := fs.List(ctx, stdpath.Join(dst, sub), &fs.ListArgs{
// 因为alias是NoCache且Get方法不会返回NotSupport或NotImplement错误
// 所以这里对象不会传回到alias,也就不需要返回BalancedObjs了
objMap := make(map[string]model.Obj)
for _, dir := range dirs {
if dir == nil {
continue
}
dirPath := dir.GetPath()
tmp, err := fs.List(ctx, dirPath, &fs.ListArgs{
NoLog: true,
Refresh: args.Refresh,
WithStorageDetails: args.WithStorageDetails && d.DetailsPassThrough,
})
if err == nil {
tmp, err = utils.SliceConvert(tmp, func(obj model.Obj) (model.Obj, error) {
objRes := model.Object{
Name: obj.GetName(),
Size: obj.GetSize(),
Modified: obj.ModTime(),
IsFolder: obj.IsDir(),
}
if thumb, ok := model.GetThumb(obj); ok {
return &model.ObjThumb{
Object: objRes,
Thumbnail: model.Thumbnail{
Thumbnail: thumb,
},
}, nil
}
if details, ok := model.GetStorageDetails(obj); ok {
return &model.ObjStorageDetails{
Obj: &objRes,
StorageDetailsWithName: *details,
}, nil
}
return &objRes, nil
})
if err != nil {
continue
}
if err == nil {
objs = append(objs, tmp...)
for _, obj := range tmp {
name := obj.GetName()
if _, exists := objMap[name]; exists {
continue
}
mask := model.GetObjMask(obj) &^ model.Temp
objRes := model.Object{
Name: name,
Path: stdpath.Join(dirPath, name),
Size: obj.GetSize(),
Modified: obj.ModTime(),
IsFolder: obj.IsDir(),
Mask: mask,
}
var objRet model.Obj
if thumb, ok := model.GetThumb(obj); ok {
objRet = &model.ObjThumb{
Object: objRes,
Thumbnail: model.Thumbnail{
Thumbnail: thumb,
},
}
} else {
objRet = &objRes
}
if details, ok := model.GetStorageDetails(obj); ok {
objRet = &model.ObjStorageDetails{
Obj: objRet,
StorageDetailsWithName: *details,
}
}
objMap[name] = objRet
}
}
objs := make([]model.Obj, 0, len(objMap))
for _, obj := range objMap {
objs = append(objs, obj)
}
return objs, nil
}
func (d *Alias) Link(ctx context.Context, file model.Obj, args model.LinkArgs) (*model.Link, error) {
root, sub := d.getRootAndPath(file.GetPath())
dsts, ok := d.pathMap[root]
if !ok {
return nil, errs.ObjectNotFound
}
// proxy || ftp,s3
if common.GetApiUrl(ctx) == "" {
args.Redirect = false
}
for _, dst := range dsts {
reqPath := stdpath.Join(dst, sub)
link, fi, err := d.link(ctx, reqPath, args)
if d.ReadConflictPolicy == AllRWP && !args.Redirect {
files, err := d.getAllObjs(ctx, file, getWriteAndPutFilterFunc(AllRWP))
if err != nil {
continue
return nil, err
}
if link == nil {
// 重定向且需要通过代理
return &model.Link{
URL: fmt.Sprintf("%s/p%s?sign=%s",
common.GetApiUrl(ctx),
utils.EncodePath(reqPath, true),
sign.Sign(reqPath)),
}, nil
linkClosers := make([]io.Closer, 0, len(files))
rrf := make([]model.RangeReaderIF, 0, len(files))
for _, f := range files {
link, fi, err := d.link(ctx, f.GetPath(), args)
if err != nil {
continue
}
if fi.GetSize() != files.GetSize() {
_ = link.Close()
continue
}
l := *link // 复制一份,避免修改到原始link
if l.ContentLength == 0 {
l.ContentLength = fi.GetSize()
}
if d.DownloadConcurrency > 0 {
l.Concurrency = d.DownloadConcurrency
}
if d.DownloadPartSize > 0 {
l.PartSize = d.DownloadPartSize * utils.KB
}
rr, err := stream.GetRangeReaderFromLink(l.ContentLength, &l)
if err != nil {
_ = link.Close()
continue
}
linkClosers = append(linkClosers, link)
rrf = append(rrf, rr)
}
rr := func(ctx context.Context, httpRange http_range.Range) (io.ReadCloser, error) {
return rrf[rand.Intn(len(rrf))].RangeRead(ctx, httpRange)
}
return &model.Link{
RangeReader: stream.RangeReaderFunc(rr),
SyncClosers: utils.NewSyncClosers(linkClosers...),
}, nil
}
resultLink := *link
resultLink.SyncClosers = utils.NewSyncClosers(link)
if args.Redirect {
return &resultLink, nil
}
if resultLink.ContentLength == 0 {
resultLink.ContentLength = fi.GetSize()
}
if d.DownloadConcurrency > 0 {
resultLink.Concurrency = d.DownloadConcurrency
}
if d.DownloadPartSize > 0 {
resultLink.PartSize = d.DownloadPartSize * utils.KB
}
reqPath := d.getBalancedPath(ctx, file)
link, fi, err := d.link(ctx, reqPath, args)
if err != nil {
return nil, err
}
if link == nil {
// 重定向且需要通过代理
return &model.Link{
URL: fmt.Sprintf("%s/p%s?sign=%s",
common.GetApiUrl(ctx),
utils.EncodePath(reqPath, true),
sign.Sign(reqPath)),
}, nil
}
resultLink := *link // 复制一份,避免修改到原始link
resultLink.SyncClosers = utils.NewSyncClosers(link)
if args.Redirect {
return &resultLink, nil
}
return nil, errs.ObjectNotFound
if resultLink.ContentLength == 0 {
resultLink.ContentLength = fi.GetSize()
}
if d.DownloadConcurrency > 0 {
resultLink.Concurrency = d.DownloadConcurrency
}
if d.DownloadPartSize > 0 {
resultLink.PartSize = d.DownloadPartSize * utils.KB
}
return &resultLink, nil
}
func (d *Alias) Other(ctx context.Context, args model.OtherArgs) (interface{}, error) {
root, sub := d.getRootAndPath(args.Obj.GetPath())
dsts, ok := d.pathMap[root]
if !ok {
return nil, errs.ObjectNotFound
// Other 不应负载均衡,这是因为前端是否调用 /fs/other 的判断条件是返回的 provider 的值
// 而 ProviderPassThrough 开启时,返回的 provider 固定为第一个 obj 的后端驱动
storage, actualPath, err := op.GetStorageAndActualPath(args.Obj.GetPath())
if err != nil {
return nil, err
}
for _, dst := range dsts {
rawPath := stdpath.Join(dst, sub)
storage, actualPath, err := op.GetStorageAndActualPath(rawPath)
if err != nil {
continue
}
other, ok := storage.(driver.Other)
if !ok {
continue
}
obj, err := op.GetUnwrap(ctx, storage, actualPath)
if err != nil {
continue
}
return other.Other(ctx, model.OtherArgs{
Obj: obj,
Method: args.Method,
Data: args.Data,
})
}
return nil, errs.NotImplement
return op.Other(ctx, storage, model.FsOtherArgs{
Path: actualPath,
Method: args.Method,
Data: args.Data,
})
}
func (d *Alias) MakeDir(ctx context.Context, parentDir model.Obj, dirName string) error {
if !d.Writable {
return errs.PermissionDenied
}
reqPath, err := d.getReqPath(ctx, parentDir, true)
objs, err := d.getWriteObjs(ctx, parentDir)
if err == nil {
for _, path := range reqPath {
err = errors.Join(err, fs.MakeDir(ctx, stdpath.Join(*path, dirName)))
for _, obj := range objs {
err = errors.Join(err, fs.MakeDir(ctx, stdpath.Join(obj.GetPath(), dirName)))
}
return err
}
if errs.IsNotImplementError(err) {
return errors.New("same-name dirs cannot make sub-dir")
}
return err
}
func (d *Alias) Move(ctx context.Context, srcObj, dstDir model.Obj) error {
if !d.Writable {
return errs.PermissionDenied
}
srcPath, err := d.getReqPath(ctx, srcObj, false)
if errs.IsNotImplementError(err) {
return errors.New("same-name files cannot be moved")
}
if err != nil {
return err
}
dstPath, err := d.getReqPath(ctx, dstDir, true)
if errs.IsNotImplementError(err) {
return errors.New("same-name dirs cannot be moved to")
}
if err != nil {
return err
}
if len(srcPath) == len(dstPath) {
for i := range srcPath {
_, e := fs.Move(ctx, *srcPath[i], *dstPath[i])
srcs, dsts, err := d.getMoveObjs(ctx, srcObj, dstDir)
if err == nil {
for i, dst := range dsts {
src := srcs[i]
_, e := fs.Move(ctx, src.GetPath(), dst.GetPath())
err = errors.Join(err, e)
}
srcs = srcs[len(dsts):]
for _, src := range srcs {
e := fs.Remove(ctx, src.GetPath())
err = errors.Join(err, e)
}
return err
} else {
return errors.New("parallel paths mismatch")
}
return err
}
func (d *Alias) Rename(ctx context.Context, srcObj model.Obj, newName string) error {
if !d.Writable {
return errs.PermissionDenied
}
reqPath, err := d.getReqPath(ctx, srcObj, false)
objs, err := d.getWriteObjs(ctx, srcObj)
if err == nil {
for _, path := range reqPath {
err = errors.Join(err, fs.Rename(ctx, *path, newName))
for _, obj := range objs {
err = errors.Join(err, fs.Rename(ctx, obj.GetPath(), newName))
}
return err
}
if errs.IsNotImplementError(err) {
return errors.New("same-name files cannot be Rename")
}
return err
}
func (d *Alias) Copy(ctx context.Context, srcObj, dstDir model.Obj) error {
if !d.Writable {
return errs.PermissionDenied
}
srcPath, err := d.getReqPath(ctx, srcObj, false)
if errs.IsNotImplementError(err) {
return errors.New("same-name files cannot be copied")
}
if err != nil {
return err
}
dstPath, err := d.getReqPath(ctx, dstDir, true)
if errs.IsNotImplementError(err) {
return errors.New("same-name dirs cannot be copied to")
}
if err != nil {
return err
}
if len(srcPath) == len(dstPath) {
for i := range srcPath {
_, e := fs.Copy(ctx, *srcPath[i], *dstPath[i])
srcs, dsts, err := d.getCopyObjs(ctx, srcObj, dstDir)
if err == nil {
for i, src := range srcs {
dst := dsts[i]
_, e := fs.Copy(ctx, src.GetPath(), dst.GetPath())
err = errors.Join(err, e)
}
return err
} else if len(srcPath) == 1 || !d.ProtectSameName {
for _, path := range dstPath {
_, e := fs.Copy(ctx, *srcPath[0], *path)
err = errors.Join(err, e)
}
return err
} else {
return errors.New("parallel paths mismatch")
}
return err
}
func (d *Alias) Remove(ctx context.Context, obj model.Obj) error {
if !d.Writable {
return errs.PermissionDenied
}
reqPath, err := d.getReqPath(ctx, obj, false)
objs, err := d.getWriteObjs(ctx, obj)
if err == nil {
for _, path := range reqPath {
err = errors.Join(err, fs.Remove(ctx, *path))
for _, obj := range objs {
err = errors.Join(err, fs.Remove(ctx, obj.GetPath()))
}
return err
}
if errs.IsNotImplementError(err) {
return errors.New("same-name files cannot be Delete")
}
return err
}
func (d *Alias) Put(ctx context.Context, dstDir model.Obj, s model.FileStreamer, up driver.UpdateProgress) error {
if !d.Writable {
return errs.PermissionDenied
}
reqPath, err := d.getReqPath(ctx, dstDir, true)
objs, err := d.getPutObjs(ctx, dstDir)
if err == nil {
if len(reqPath) == 1 {
storage, reqActualPath, err := op.GetStorageAndActualPath(*reqPath[0])
if len(objs) == 1 {
storage, reqActualPath, err := op.GetStorageAndActualPath(objs.GetPath())
if err != nil {
return err
}
@@ -387,10 +398,10 @@ func (d *Alias) Put(ctx context.Context, dstDir model.Obj, s model.FileStreamer,
if err != nil {
return err
}
count := float64(len(reqPath) + 1)
count := float64(len(objs) + 1)
up(100 / count)
for i, path := range reqPath {
err = errors.Join(err, fs.PutDirectly(ctx, *path, &stream.FileStream{
for i, obj := range objs {
err = errors.Join(err, fs.PutDirectly(ctx, obj.GetPath(), &stream.FileStream{
Obj: s,
Mimetype: s.GetMimetype(),
Reader: file,
@@ -404,55 +415,40 @@ func (d *Alias) Put(ctx context.Context, dstDir model.Obj, s model.FileStreamer,
return err
}
}
if errs.IsNotImplementError(err) {
return errors.New("same-name dirs cannot be Put")
}
return err
}
func (d *Alias) PutURL(ctx context.Context, dstDir model.Obj, name, url string) error {
if !d.Writable {
return errs.PermissionDenied
}
reqPath, err := d.getReqPath(ctx, dstDir, true)
objs, err := d.getPutObjs(ctx, dstDir)
if err == nil {
for _, path := range reqPath {
err = errors.Join(err, fs.PutURL(ctx, *path, name, url))
for _, obj := range objs {
err = errors.Join(err, fs.PutURL(ctx, obj.GetPath(), name, url))
}
return err
}
if errs.IsNotImplementError(err) {
return errors.New("same-name files cannot offline download")
}
return err
}
func (d *Alias) GetArchiveMeta(ctx context.Context, obj model.Obj, args model.ArchiveArgs) (model.ArchiveMeta, error) {
root, sub := d.getRootAndPath(obj.GetPath())
dsts, ok := d.pathMap[root]
if !ok {
return nil, errs.ObjectNotFound
reqPath := d.getBalancedPath(ctx, obj)
if reqPath == "" {
return nil, errs.NotFile
}
for _, dst := range dsts {
meta, err := d.getArchiveMeta(ctx, dst, sub, args)
if err == nil {
return meta, nil
}
meta, err := d.getArchiveMeta(ctx, reqPath, args)
if err == nil {
return meta, nil
}
return nil, errs.NotImplement
}
func (d *Alias) ListArchive(ctx context.Context, obj model.Obj, args model.ArchiveInnerArgs) ([]model.Obj, error) {
root, sub := d.getRootAndPath(obj.GetPath())
dsts, ok := d.pathMap[root]
if !ok {
return nil, errs.ObjectNotFound
reqPath := d.getBalancedPath(ctx, obj)
if reqPath == "" {
return nil, errs.NotFile
}
for _, dst := range dsts {
l, err := d.listArchive(ctx, dst, sub, args)
if err == nil {
return l, nil
}
l, err := d.listArchive(ctx, reqPath, args)
if err == nil {
return l, nil
}
return nil, errs.NotImplement
}
@@ -461,67 +457,62 @@ func (d *Alias) Extract(ctx context.Context, obj model.Obj, args model.ArchiveIn
// alias的两个驱动,一个支持驱动提取,一个不支持,如何兼容?
// 如果访问的是不支持驱动提取的驱动内的压缩文件,GetArchiveMeta就会返回errs.NotImplement,提取URL前缀就会是/ae,Extract就不会被调用
// 如果访问的是支持驱动提取的驱动内的压缩文件,GetArchiveMeta就会返回有效值,提取URL前缀就会是/ad,Extract就会被调用
root, sub := d.getRootAndPath(obj.GetPath())
dsts, ok := d.pathMap[root]
if !ok {
return nil, errs.ObjectNotFound
reqPath := d.getBalancedPath(ctx, obj)
if reqPath == "" {
return nil, errs.NotFile
}
for _, dst := range dsts {
reqPath := stdpath.Join(dst, sub)
link, err := d.extract(ctx, reqPath, args)
if err != nil {
continue
}
if link == nil {
return &model.Link{
URL: fmt.Sprintf("%s/ap%s?inner=%s&pass=%s&sign=%s",
common.GetApiUrl(ctx),
utils.EncodePath(reqPath, true),
utils.EncodePath(args.InnerPath, true),
url.QueryEscape(args.Password),
sign.SignArchive(reqPath)),
}, nil
}
resultLink := *link
resultLink.SyncClosers = utils.NewSyncClosers(link)
return &resultLink, nil
link, err := d.extract(ctx, reqPath, args)
if err != nil {
return nil, errs.NotImplement
}
return nil, errs.NotImplement
if link == nil {
return &model.Link{
URL: fmt.Sprintf("%s/ap%s?inner=%s&pass=%s&sign=%s",
common.GetApiUrl(ctx),
utils.EncodePath(reqPath, true),
utils.EncodePath(args.InnerPath, true),
url.QueryEscape(args.Password),
sign.SignArchive(reqPath)),
}, nil
}
resultLink := *link
resultLink.SyncClosers = utils.NewSyncClosers(link)
return &resultLink, nil
}
func (d *Alias) ArchiveDecompress(ctx context.Context, srcObj, dstDir model.Obj, args model.ArchiveDecompressArgs) error {
if !d.Writable {
return errs.PermissionDenied
}
srcPath, err := d.getReqPath(ctx, srcObj, false)
if errs.IsNotImplementError(err) {
return errors.New("same-name files cannot be decompressed")
}
if err != nil {
return err
}
dstPath, err := d.getReqPath(ctx, dstDir, true)
if errs.IsNotImplementError(err) {
return errors.New("same-name dirs cannot be decompressed to")
}
if err != nil {
return err
}
if len(srcPath) == len(dstPath) {
for i := range srcPath {
_, e := fs.ArchiveDecompress(ctx, *srcPath[i], *dstPath[i], args)
srcs, dsts, err := d.getCopyObjs(ctx, srcObj, dstDir)
if err == nil {
for i, src := range srcs {
dst := dsts[i]
_, e := fs.ArchiveDecompress(ctx, src.GetPath(), dst.GetPath(), args)
err = errors.Join(err, e)
}
return err
} else if len(srcPath) == 1 || !d.ProtectSameName {
for _, path := range dstPath {
_, e := fs.ArchiveDecompress(ctx, *srcPath[0], *path, args)
err = errors.Join(err, e)
}
return err
} else {
return errors.New("parallel paths mismatch")
}
return err
}
func (d *Alias) ResolveLinkCacheMode(path string) driver.LinkCacheMode {
roots, sub := d.getRootsAndPath(path)
if len(roots) == 0 {
return 0
}
for _, root := range roots {
storage, actualPath, err := op.GetStorageAndActualPath(stdpath.Join(root, sub))
if err != nil {
continue
}
if storage.Config().CheckStatus && storage.GetStorage().Status != op.WORK {
continue
}
mode := storage.Config().LinkCacheMode
if mode == -1 {
return storage.(driver.LinkCacheModeResolver).ResolveLinkCacheMode(actualPath)
} else {
return mode
}
}
return 0
}
var _ driver.Driver = (*Alias)(nil)
+11 -16
View File
@@ -6,17 +6,15 @@ import (
)
type Addition struct {
// Usually one of two
// driver.RootPath
// define other
Paths string `json:"paths" required:"true" type:"text"`
ProtectSameName bool `json:"protect_same_name" default:"true" required:"false" help:"Protects same-name files from Delete or Rename"`
ParallelWrite bool `json:"parallel_write" type:"bool" default:"false"`
DownloadConcurrency int `json:"download_concurrency" default:"0" required:"false" type:"number" help:"Need to enable proxy"`
DownloadPartSize int `json:"download_part_size" default:"0" type:"number" required:"false" help:"Need to enable proxy. Unit: KB"`
Writable bool `json:"writable" type:"bool" default:"false"`
ProviderPassThrough bool `json:"provider_pass_through" type:"bool" default:"false"`
DetailsPassThrough bool `json:"details_pass_through" type:"bool" default:"false"`
Paths string `json:"paths" required:"true" type:"text"`
ReadConflictPolicy string `json:"read_conflict_policy" type:"select" options:"first,random,all" default:"first"`
WriteConflictPolicy string `json:"write_conflict_policy" type:"select" options:"disabled,first,deterministic,deterministic_or_all,all,all_strict" default:"disabled" help:"How the driver handles identical backend paths when renaming, removing, or making directories."`
PutConflictPolicy string `json:"put_conflict_policy" type:"select" options:"disabled,first,deterministic,deterministic_or_all,all,all_strict,random,quota,quota_strict" default:"disabled" help:"How the driver handles identical backend paths when uploading, copying, moving, or decompressing."`
FileConsistencyCheck bool `json:"file_consistency_check" type:"bool" default:"false"`
DownloadConcurrency int `json:"download_concurrency" default:"0" required:"false" type:"number" help:"Need to enable proxy"`
DownloadPartSize int `json:"download_part_size" default:"0" type:"number" required:"false" help:"Need to enable proxy. Unit: KB"`
ProviderPassThrough bool `json:"provider_pass_through" type:"bool" default:"false"`
DetailsPassThrough bool `json:"details_pass_through" type:"bool" default:"false"`
}
var config = driver.Config{
@@ -26,14 +24,11 @@ var config = driver.Config{
NoUpload: false,
DefaultRoot: "/",
ProxyRangeOption: true,
LinkCacheMode: driver.LinkCacheAuto,
}
func init() {
op.RegisterDriver(func() driver.Driver {
return &Alias{
Addition: Addition{
ProtectSameName: true,
},
}
return &Alias{}
})
}
+77
View File
@@ -1 +1,78 @@
package alias
import (
"time"
"github.com/OpenListTeam/OpenList/v4/internal/model"
"github.com/OpenListTeam/OpenList/v4/pkg/utils"
"github.com/pkg/errors"
)
const (
DisabledWP = "disabled"
FirstRWP = "first"
DeterministicWP = "deterministic"
DeterministicOrAllWP = "deterministic_or_all"
AllRWP = "all"
AllStrictWP = "all_strict"
RandomBalancedRP = "random"
BalancedByQuotaP = "quota"
BalancedByQuotaStrictP = "quota_strict"
)
var (
ValidReadConflictPolicy = []string{FirstRWP, RandomBalancedRP, AllRWP}
ValidWriteConflictPolicy = []string{DisabledWP, FirstRWP, DeterministicWP, DeterministicOrAllWP, AllRWP,
AllStrictWP}
ValidPutConflictPolicy = []string{DisabledWP, FirstRWP, DeterministicWP, DeterministicOrAllWP, AllRWP,
AllStrictWP, RandomBalancedRP, BalancedByQuotaP, BalancedByQuotaStrictP}
)
var (
ErrPathConflict = errors.New("path conflict")
ErrSamePathLeak = errors.New("leak some of same-name dirs")
ErrNoEnoughSpace = errors.New("none of same-name dirs has enough space")
ErrNotEnoughSrcObjs = errors.New("cannot move fewer objs to more paths, please try copying")
)
type BalancedObjs []model.Obj
func (b BalancedObjs) GetSize() int64 {
return b[0].GetSize()
}
func (b BalancedObjs) ModTime() time.Time {
return b[0].ModTime()
}
func (b BalancedObjs) CreateTime() time.Time {
return b[0].CreateTime()
}
func (b BalancedObjs) IsDir() bool {
return b[0].IsDir()
}
func (b BalancedObjs) GetHash() utils.HashInfo {
return b[0].GetHash()
}
func (b BalancedObjs) GetName() string {
return b[0].GetName()
}
func (b BalancedObjs) GetPath() string {
return b[0].GetPath()
}
func (b BalancedObjs) GetID() string {
return b[0].GetID()
}
func (b BalancedObjs) Unwrap() model.Obj {
return b[0]
}
var _ model.Obj = (BalancedObjs)(nil)
type tempObj struct{ model.Object }
+374 -75
View File
@@ -2,10 +2,9 @@ package alias
import (
"context"
"errors"
"math/rand"
stdpath "path"
"strings"
"sync"
"time"
"github.com/OpenListTeam/OpenList/v4/internal/driver"
@@ -14,20 +13,29 @@ import (
"github.com/OpenListTeam/OpenList/v4/internal/model"
"github.com/OpenListTeam/OpenList/v4/internal/op"
"github.com/OpenListTeam/OpenList/v4/server/common"
"github.com/pkg/errors"
log "github.com/sirupsen/logrus"
)
type detailWithIndex struct {
idx int
val *model.StorageDetails
}
func (d *Alias) listRoot(ctx context.Context, withDetails, refresh bool) []model.Obj {
var objs []model.Obj
var wg sync.WaitGroup
detailsChan := make(chan detailWithIndex, len(d.pathMap))
workerCount := 0
for _, k := range d.rootOrder {
obj := model.Object{
obj := &model.Object{
Name: k,
Path: "/" + k,
IsFolder: true,
Modified: d.Modified,
Mask: model.Locked | model.Virtual,
}
idx := len(objs)
objs = append(objs, &obj)
objs = append(objs, obj)
v := d.pathMap[k]
if !withDetails || len(v) != 1 {
continue
@@ -36,6 +44,7 @@ func (d *Alias) listRoot(ctx context.Context, withDetails, refresh bool) []model
if err != nil {
continue
}
obj.Modified = remoteDriver.GetStorage().Modified
_, ok := remoteDriver.(driver.WithDetails)
if !ok {
continue
@@ -47,47 +56,47 @@ func (d *Alias) listRoot(ctx context.Context, withDetails, refresh bool) []model
DriverName: remoteDriver.Config().Name,
},
}
wg.Add(1)
go func() {
defer wg.Done()
c, cancel := context.WithTimeout(ctx, time.Second)
defer cancel()
details, e := op.GetStorageDetails(c, remoteDriver, refresh)
workerCount++
go func(dri driver.Driver, i int) {
details, e := op.GetStorageDetails(ctx, dri, refresh)
if e != nil {
if !errors.Is(e, errs.NotImplement) && !errors.Is(e, errs.StorageNotInit) {
log.Errorf("failed get %s storage details: %+v", remoteDriver.GetStorage().MountPath, e)
log.Errorf("failed get %s storage details: %+v", dri.GetStorage().MountPath, e)
}
return
}
objs[idx].(*model.ObjStorageDetails).StorageDetails = details
}()
detailsChan <- detailWithIndex{idx: i, val: details}
}(remoteDriver, idx)
}
for workerCount > 0 {
select {
case r := <-detailsChan:
objs[r.idx].(*model.ObjStorageDetails).StorageDetails = r.val
workerCount--
case <-time.After(time.Second):
workerCount = 0
}
}
wg.Wait()
return objs
}
// do others that not defined in Driver interface
func getPair(path string) (string, string) {
// path = strings.TrimSpace(path)
if strings.Contains(path, ":") {
pair := strings.SplitN(path, ":", 2)
if !strings.Contains(pair[0], "/") {
return pair[0], pair[1]
}
if name, path, ok := strings.Cut(path, ":"); ok && !strings.Contains(name, "/") {
return name, path
}
return stdpath.Base(path), path
}
func (d *Alias) getRootAndPath(path string) (string, string) {
if d.autoFlatten {
return d.oneKey, path
func (d *Alias) getRootsAndPath(path string) (roots []string, sub string) {
if len(d.rootOrder) == 1 {
return d.pathMap[d.rootOrder[0]], path
}
path = strings.TrimPrefix(path, "/")
parts := strings.SplitN(path, "/", 2)
if len(parts) == 1 {
return parts[0], ""
before, after, ok := strings.Cut(path, "/")
if !ok {
return d.pathMap[path], ""
}
return parts[0], parts[1]
return d.pathMap[before], after
}
func (d *Alias) link(ctx context.Context, reqPath string, args model.LinkArgs) (*model.Link, model.Obj, error) {
@@ -95,59 +104,350 @@ func (d *Alias) link(ctx context.Context, reqPath string, args model.LinkArgs) (
if err != nil {
return nil, nil, err
}
if !args.Redirect {
return op.Link(ctx, storage, reqActualPath, args)
}
obj, err := fs.Get(ctx, reqPath, &fs.GetArgs{NoLog: true})
if err != nil {
return nil, nil, err
}
if common.ShouldProxy(storage, stdpath.Base(reqPath)) {
return nil, obj, nil
if args.Redirect && common.ShouldProxy(storage, stdpath.Base(reqPath)) {
return nil, nil, nil
}
return op.Link(ctx, storage, reqActualPath, args)
}
func (d *Alias) getReqPath(ctx context.Context, obj model.Obj, isParent bool) ([]*string, error) {
root, sub := d.getRootAndPath(obj.GetPath())
if sub == "" && !isParent {
return nil, errs.NotSupport
func isConsistent(a, b model.Obj) bool {
if a.GetSize() != b.GetSize() {
return false
}
dsts, ok := d.pathMap[root]
all := true
if !ok {
return nil, errs.ObjectNotFound
}
var reqPath []*string
for _, dst := range dsts {
path := stdpath.Join(dst, sub)
_, err := fs.Get(ctx, path, &fs.GetArgs{NoLog: true})
if err != nil {
all = false
if d.ProtectSameName && d.ParallelWrite && len(reqPath) >= 2 {
return nil, errs.NotImplement
}
continue
}
if !d.ProtectSameName && !d.ParallelWrite {
return []*string{&path}, nil
}
reqPath = append(reqPath, &path)
if d.ProtectSameName && !d.ParallelWrite && len(reqPath) >= 2 {
return nil, errs.NotImplement
}
if d.ProtectSameName && d.ParallelWrite && len(reqPath) >= 2 && !all {
return nil, errs.NotImplement
for ht, v := range a.GetHash().All() {
ah := b.GetHash().GetHash(ht)
if ah != "" && ah != v {
return false
}
}
if len(reqPath) == 0 {
return nil, errs.ObjectNotFound
}
return reqPath, nil
return true
}
func (d *Alias) getArchiveMeta(ctx context.Context, dst, sub string, args model.ArchiveArgs) (model.ArchiveMeta, error) {
reqPath := stdpath.Join(dst, sub)
func (d *Alias) getAllObjs(ctx context.Context, bObj model.Obj, ifContinue func(err error) (bool, error)) (BalancedObjs, error) {
objs := bObj.(BalancedObjs)
length := 0
for _, o := range objs {
var err error
var obj model.Obj
temp, isTemp := o.(*tempObj)
if isTemp {
obj, err = fs.Get(ctx, o.GetPath(), &fs.GetArgs{NoLog: true})
if err == nil {
if !bObj.IsDir() {
if obj.IsDir() {
err = errs.NotFile
} else if d.FileConsistencyCheck && !isConsistent(bObj, obj) {
err = errs.ObjectNotFound
}
} else if !obj.IsDir() {
err = errs.NotFolder
}
}
} else if o == nil {
err = errs.ObjectNotFound
}
cont, err := ifContinue(err)
if err != nil {
if cont {
continue
}
return nil, err
}
if isTemp {
objRes := temp.Object
// objRes.Name = obj.GetName()
// objRes.Size = obj.GetSize()
// objRes.Modified = obj.ModTime()
// objRes.HashInfo = obj.GetHash()
objs[length] = &objRes
} else {
objs[length] = o
}
length++
if !cont {
break
}
}
if length == 0 {
return nil, errs.ObjectNotFound
}
return objs[:length], nil
}
func (d *Alias) getBalancedPath(ctx context.Context, file model.Obj) string {
if d.ReadConflictPolicy == FirstRWP {
return file.GetPath()
}
files := file.(BalancedObjs)
if rand.Intn(len(files)) == 0 {
return file.GetPath()
}
files, _ = d.getAllObjs(ctx, file, getWriteAndPutFilterFunc(AllRWP))
return files[rand.Intn(len(files))].GetPath()
}
func getWriteAndPutFilterFunc(policy string) func(error) (bool, error) {
if policy == AllRWP {
return func(err error) (bool, error) {
return true, err
}
}
all := true
l := 0
return func(err error) (bool, error) {
if err != nil {
switch policy {
case AllStrictWP:
return false, ErrSamePathLeak
case DeterministicOrAllWP:
if l >= 2 {
return false, ErrSamePathLeak
}
}
all = false
} else {
switch policy {
case FirstRWP:
return false, nil
case DeterministicWP:
if l > 0 {
return false, ErrPathConflict
}
case DeterministicOrAllWP:
if l > 0 && !all {
return false, ErrSamePathLeak
}
}
l += 1
}
return true, err
}
}
func (d *Alias) getWriteObjs(ctx context.Context, obj model.Obj) (BalancedObjs, error) {
if d.WriteConflictPolicy == DisabledWP {
return nil, errs.PermissionDenied
}
return d.getAllObjs(ctx, obj, getWriteAndPutFilterFunc(d.WriteConflictPolicy))
}
func (d *Alias) getPutObjs(ctx context.Context, obj model.Obj) (BalancedObjs, error) {
if d.PutConflictPolicy == DisabledWP {
return nil, errs.PermissionDenied
}
objs, err := d.getAllObjs(ctx, obj, getWriteAndPutFilterFunc(d.PutConflictPolicy))
if err != nil {
return nil, err
}
strict := false
switch d.PutConflictPolicy {
case RandomBalancedRP:
ri := rand.Intn(len(objs))
return objs[ri : ri+1], nil
case BalancedByQuotaStrictP:
strict = true
fallthrough
case BalancedByQuotaP:
objs, ok := getRandomObjByQuotaBalanced(ctx, objs, strict, uint64(obj.GetSize()))
if !ok {
return nil, ErrNoEnoughSpace
}
return objs, nil
default:
return objs, nil
}
}
func getRandomObjByQuotaBalanced(ctx context.Context, reqPath BalancedObjs, strict bool, objSize uint64) (BalancedObjs, bool) {
// Get all space
details := make([]*model.StorageDetails, len(reqPath))
detailsChan := make(chan detailWithIndex, len(reqPath))
workerCount := 0
for i, p := range reqPath {
s, err := fs.GetStorage(p.GetPath(), &fs.GetStoragesArgs{})
if err != nil {
continue
}
if _, ok := s.(driver.WithDetails); !ok {
continue
}
workerCount++
go func(dri driver.Driver, i int) {
d, e := op.GetStorageDetails(ctx, dri)
if e != nil {
if !errors.Is(e, errs.NotImplement) && !errors.Is(e, errs.StorageNotInit) {
log.Errorf("failed get %s storage details: %+v", dri.GetStorage().MountPath, e)
}
}
detailsChan <- detailWithIndex{idx: i, val: d}
}(s, i)
}
for workerCount > 0 {
select {
case r := <-detailsChan:
details[r.idx] = r.val
workerCount--
case <-time.After(time.Second):
workerCount = 0
}
}
// Try select one that has space info
selected, ok := selectRandom(details, func(d *model.StorageDetails) uint64 {
if d == nil || d.FreeSpace < objSize {
return 0
}
return d.FreeSpace
})
if !ok {
if strict {
return nil, false
} else {
// No strict mode, return any of non-details ones
noDetails := make([]int, 0, len(details))
for i, d := range details {
if d == nil {
noDetails = append(noDetails, i)
}
}
if len(noDetails) == 0 {
return nil, false
}
selected = noDetails[rand.Intn(len(noDetails))]
}
}
return reqPath[selected : selected+1], true
}
func selectRandom[Item any](arr []Item, getWeight func(Item) uint64) (int, bool) {
var totalWeight uint64 = 0
for _, i := range arr {
totalWeight += getWeight(i)
}
if totalWeight == 0 {
return 0, false
}
r := rand.Uint64() % totalWeight
for i, item := range arr {
w := getWeight(item)
if r < w {
return i, true
}
r -= w
}
return 0, false
}
func (d *Alias) getCopyObjs(ctx context.Context, srcObj, dstDir model.Obj) (BalancedObjs, BalancedObjs, error) {
if d.PutConflictPolicy == DisabledWP {
return nil, nil, errs.PermissionDenied
}
dstObjs, err := d.getAllObjs(ctx, dstDir, getWriteAndPutFilterFunc(d.PutConflictPolicy))
if err != nil {
return nil, nil, err
}
dstStorageMap := make(map[string][]model.Obj)
allocatingDst := make(map[model.Obj]struct{})
for _, o := range dstObjs {
storage, e := fs.GetStorage(o.GetPath(), &fs.GetStoragesArgs{})
if e != nil {
return nil, nil, errors.WithMessagef(e, "cannot copy to virtual path [%s]", o.GetPath())
}
mp := storage.GetStorage().MountPath
dstStorageMap[mp] = append(dstStorageMap[mp], o)
allocatingDst[o] = struct{}{}
}
tmpSrcObjs, err := d.getAllObjs(ctx, srcObj, getWriteAndPutFilterFunc(AllRWP))
if err != nil {
return nil, nil, err
}
srcObjs := make(BalancedObjs, 0, len(dstObjs))
for _, src := range tmpSrcObjs {
storage, e := fs.GetStorage(src.GetPath(), &fs.GetStoragesArgs{})
if e != nil {
continue
}
mp := storage.GetStorage().MountPath
if tmp, ok := dstStorageMap[mp]; ok {
for _, dst := range tmp {
dstObjs[len(srcObjs)] = dst
srcObjs = append(srcObjs, src)
delete(allocatingDst, dst)
}
delete(dstStorageMap, mp)
}
}
dstObjs = dstObjs[:len(srcObjs)]
for dst := range allocatingDst {
src := tmpSrcObjs[0]
if d.ReadConflictPolicy == RandomBalancedRP || d.ReadConflictPolicy == AllRWP {
src = tmpSrcObjs[rand.Intn(len(tmpSrcObjs))]
}
srcObjs = append(srcObjs, src)
dstObjs = append(dstObjs, dst)
}
return srcObjs, dstObjs, nil
}
func (d *Alias) getMoveObjs(ctx context.Context, srcObj, dstDir model.Obj) (BalancedObjs, BalancedObjs, error) {
if d.PutConflictPolicy == DisabledWP {
return nil, nil, errs.PermissionDenied
}
dstObjs, err := d.getAllObjs(ctx, dstDir, getWriteAndPutFilterFunc(d.PutConflictPolicy))
if err != nil {
return nil, nil, err
}
tmpSrcObjs, err := d.getAllObjs(ctx, srcObj, getWriteAndPutFilterFunc(AllRWP))
if err != nil {
return nil, nil, err
}
if len(tmpSrcObjs) < len(dstObjs) {
return nil, nil, ErrNotEnoughSrcObjs
}
dstStorageMap := make(map[string][]model.Obj)
allocatingDst := make(map[model.Obj]struct{})
for _, o := range dstObjs {
storage, e := fs.GetStorage(o.GetPath(), &fs.GetStoragesArgs{})
if e != nil {
return nil, nil, errors.WithMessagef(e, "cannot move to virtual path [%s]", o.GetPath())
}
mp := storage.GetStorage().MountPath
dstStorageMap[mp] = append(dstStorageMap[mp], o)
allocatingDst[o] = struct{}{}
}
srcObjs := make(BalancedObjs, 0, len(tmpSrcObjs))
restSrcObjs := make(BalancedObjs, 0, len(tmpSrcObjs)-len(dstObjs))
for _, src := range tmpSrcObjs {
storage, e := fs.GetStorage(src.GetPath(), &fs.GetStoragesArgs{})
if e != nil {
continue
}
mp := storage.GetStorage().MountPath
if tmp, ok := dstStorageMap[mp]; ok {
dst := tmp[0]
if len(tmp) == 1 {
delete(dstStorageMap, mp)
} else {
dstStorageMap[mp] = tmp[1:]
}
dstObjs[len(srcObjs)] = dst
srcObjs = append(srcObjs, src)
delete(allocatingDst, dst)
} else {
restSrcObjs = append(restSrcObjs, src)
}
}
dstObjs = dstObjs[:len(srcObjs)]
// len(restSrcObjs) >= len(allocatingDst)
srcObjs = append(srcObjs, restSrcObjs...)
for dst := range allocatingDst {
dstObjs = append(dstObjs, dst)
}
return srcObjs, dstObjs, nil
}
func (d *Alias) getArchiveMeta(ctx context.Context, reqPath string, args model.ArchiveArgs) (model.ArchiveMeta, error) {
storage, reqActualPath, err := op.GetStorageAndActualPath(reqPath)
if err != nil {
return nil, err
@@ -161,8 +461,7 @@ func (d *Alias) getArchiveMeta(ctx context.Context, dst, sub string, args model.
return nil, errs.NotImplement
}
func (d *Alias) listArchive(ctx context.Context, dst, sub string, args model.ArchiveInnerArgs) ([]model.Obj, error) {
reqPath := stdpath.Join(dst, sub)
func (d *Alias) listArchive(ctx context.Context, reqPath string, args model.ArchiveInnerArgs) ([]model.Obj, error) {
storage, reqActualPath, err := op.GetStorageAndActualPath(reqPath)
if err != nil {
return nil, err
+379
View File
@@ -0,0 +1,379 @@
package alist_v3
import (
"context"
"fmt"
"io"
"net/http"
"net/url"
"path"
"strings"
"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/pkg/utils"
"github.com/OpenListTeam/OpenList/v4/server/common"
"github.com/go-resty/resty/v2"
log "github.com/sirupsen/logrus"
)
type AListV3 struct {
model.Storage
Addition
}
func (d *AListV3) Config() driver.Config {
return config
}
func (d *AListV3) GetAddition() driver.Additional {
return &d.Addition
}
func (d *AListV3) Init(ctx context.Context) error {
d.Addition.Address = strings.TrimSuffix(d.Addition.Address, "/")
var resp common.Resp[MeResp]
_, _, err := d.request("/me", http.MethodGet, func(req *resty.Request) {
req.SetResult(&resp)
})
if err != nil {
return err
}
// if the username is not empty and the username is not the same as the current username, then login again
if d.Username != resp.Data.Username {
err = d.login()
if err != nil {
return err
}
}
// re-get the user info
_, _, err = d.request("/me", http.MethodGet, func(req *resty.Request) {
req.SetResult(&resp)
})
if err != nil {
return err
}
if utils.SliceContains(resp.Data.Role, model.GUEST) {
u := d.Address + "/api/public/settings"
res, err := base.RestyClient.R().Get(u)
if err != nil {
return err
}
allowMounted := utils.Json.Get(res.Body(), "data", conf.AllowMounted).ToString() == "true"
if !allowMounted {
return fmt.Errorf("the site does not allow mounted")
}
}
return err
}
func (d *AListV3) Drop(ctx context.Context) error {
return nil
}
func (d *AListV3) List(ctx context.Context, dir model.Obj, args model.ListArgs) ([]model.Obj, error) {
var resp common.Resp[FsListResp]
_, _, err := d.request("/fs/list", http.MethodPost, func(req *resty.Request) {
req.SetResult(&resp).SetBody(ListReq{
PageReq: model.PageReq{
Page: 1,
PerPage: 0,
},
Path: dir.GetPath(),
Password: d.MetaPassword,
Refresh: false,
})
})
if err != nil {
return nil, err
}
var files []model.Obj
for _, f := range resp.Data.Content {
file := model.ObjThumb{
Object: model.Object{
Name: f.Name,
Modified: f.Modified,
Ctime: f.Created,
Size: f.Size,
IsFolder: f.IsDir,
HashInfo: utils.FromString(f.HashInfo),
},
Thumbnail: model.Thumbnail{Thumbnail: f.Thumb},
}
files = append(files, &file)
}
return files, nil
}
func (d *AListV3) Link(ctx context.Context, file model.Obj, args model.LinkArgs) (*model.Link, error) {
var resp common.Resp[FsGetResp]
headers := map[string]string{
"User-Agent": base.UserAgent,
}
// if PassUAToUpsteam is true, then pass the user-agent to the upstream
if d.PassUAToUpsteam {
userAgent := args.Header.Get("user-agent")
if userAgent != "" {
headers["User-Agent"] = userAgent
}
}
// if PassIPToUpsteam is true, then pass the ip address to the upstream
if d.PassIPToUpsteam {
ip := args.IP
if ip != "" {
headers["X-Forwarded-For"] = ip
headers["X-Real-Ip"] = ip
}
}
_, _, err := d.request("/fs/get", http.MethodPost, func(req *resty.Request) {
req.SetResult(&resp).SetBody(FsGetReq{
Path: file.GetPath(),
Password: d.MetaPassword,
}).SetHeaders(headers)
})
if err != nil {
return nil, err
}
return &model.Link{
URL: resp.Data.RawURL,
}, nil
}
func (d *AListV3) MakeDir(ctx context.Context, parentDir model.Obj, dirName string) error {
_, _, err := d.request("/fs/mkdir", http.MethodPost, func(req *resty.Request) {
req.SetBody(MkdirOrLinkReq{
Path: path.Join(parentDir.GetPath(), dirName),
})
})
return err
}
func (d *AListV3) Move(ctx context.Context, srcObj, dstDir model.Obj) error {
_, _, err := d.request("/fs/move", http.MethodPost, func(req *resty.Request) {
req.SetBody(MoveCopyReq{
SrcDir: path.Dir(srcObj.GetPath()),
DstDir: dstDir.GetPath(),
Names: []string{srcObj.GetName()},
})
})
return err
}
func (d *AListV3) Rename(ctx context.Context, srcObj model.Obj, newName string) error {
_, _, err := d.request("/fs/rename", http.MethodPost, func(req *resty.Request) {
req.SetBody(RenameReq{
Path: srcObj.GetPath(),
Name: newName,
})
})
return err
}
func (d *AListV3) Copy(ctx context.Context, srcObj, dstDir model.Obj) error {
_, _, err := d.request("/fs/copy", http.MethodPost, func(req *resty.Request) {
req.SetBody(MoveCopyReq{
SrcDir: path.Dir(srcObj.GetPath()),
DstDir: dstDir.GetPath(),
Names: []string{srcObj.GetName()},
})
})
return err
}
func (d *AListV3) Remove(ctx context.Context, obj model.Obj) error {
_, _, err := d.request("/fs/remove", http.MethodPost, func(req *resty.Request) {
req.SetBody(RemoveReq{
Dir: path.Dir(obj.GetPath()),
Names: []string{obj.GetName()},
})
})
return err
}
func (d *AListV3) Put(ctx context.Context, dstDir model.Obj, s model.FileStreamer, up driver.UpdateProgress) error {
reader := driver.NewLimitedUploadStream(ctx, &driver.ReaderUpdatingProgress{
Reader: s,
UpdateProgress: up,
})
req, err := http.NewRequestWithContext(ctx, http.MethodPut, d.Address+"/api/fs/put", reader)
if err != nil {
return err
}
req.Header.Set("Authorization", d.Token)
req.Header.Set("File-Path", path.Join(dstDir.GetPath(), s.GetName()))
req.Header.Set("Password", d.MetaPassword)
if md5 := s.GetHash().GetHash(utils.MD5); len(md5) > 0 {
req.Header.Set("X-File-Md5", md5)
}
if sha1 := s.GetHash().GetHash(utils.SHA1); len(sha1) > 0 {
req.Header.Set("X-File-Sha1", sha1)
}
if sha256 := s.GetHash().GetHash(utils.SHA256); len(sha256) > 0 {
req.Header.Set("X-File-Sha256", sha256)
}
req.ContentLength = s.GetSize()
// client := base.NewHttpClient()
// client.Timeout = time.Hour * 6
res, err := base.HttpClient.Do(req)
if err != nil {
return err
}
bytes, err := io.ReadAll(res.Body)
if err != nil {
return err
}
log.Debugf("[openlist] response body: %s", string(bytes))
if res.StatusCode >= 400 {
return fmt.Errorf("request failed, status: %s", res.Status)
}
code := utils.Json.Get(bytes, "code").ToInt()
if code != 200 {
if code == 401 || code == 403 {
err = d.login()
if err != nil {
return err
}
}
return fmt.Errorf("request failed,code: %d, message: %s", code, utils.Json.Get(bytes, "message").ToString())
}
return nil
}
func (d *AListV3) GetArchiveMeta(ctx context.Context, obj model.Obj, args model.ArchiveArgs) (model.ArchiveMeta, error) {
if !d.ForwardArchiveReq {
return nil, errs.NotImplement
}
var resp common.Resp[ArchiveMetaResp]
_, code, err := d.request("/fs/archive/meta", http.MethodPost, func(req *resty.Request) {
req.SetResult(&resp).SetBody(ArchiveMetaReq{
ArchivePass: args.Password,
Password: d.MetaPassword,
Path: obj.GetPath(),
Refresh: false,
})
})
if code == 202 {
return nil, errs.WrongArchivePassword
}
if err != nil {
return nil, err
}
var tree []model.ObjTree
if resp.Data.Content != nil {
tree = make([]model.ObjTree, 0, len(resp.Data.Content))
for _, content := range resp.Data.Content {
tree = append(tree, &content)
}
}
return &model.ArchiveMetaInfo{
Comment: resp.Data.Comment,
Encrypted: resp.Data.Encrypted,
Tree: tree,
}, nil
}
func (d *AListV3) ListArchive(ctx context.Context, obj model.Obj, args model.ArchiveInnerArgs) ([]model.Obj, error) {
if !d.ForwardArchiveReq {
return nil, errs.NotImplement
}
var resp common.Resp[ArchiveListResp]
_, code, err := d.request("/fs/archive/list", http.MethodPost, func(req *resty.Request) {
req.SetResult(&resp).SetBody(ArchiveListReq{
ArchiveMetaReq: ArchiveMetaReq{
ArchivePass: args.Password,
Password: d.MetaPassword,
Path: obj.GetPath(),
Refresh: false,
},
PageReq: model.PageReq{
Page: 1,
PerPage: 0,
},
InnerPath: args.InnerPath,
})
})
if code == 202 {
return nil, errs.WrongArchivePassword
}
if err != nil {
return nil, err
}
var files []model.Obj
for _, f := range resp.Data.Content {
file := model.ObjThumb{
Object: model.Object{
Name: f.Name,
Modified: f.Modified,
Ctime: f.Created,
Size: f.Size,
IsFolder: f.IsDir,
HashInfo: utils.FromString(f.HashInfo),
},
Thumbnail: model.Thumbnail{Thumbnail: f.Thumb},
}
files = append(files, &file)
}
return files, nil
}
func (d *AListV3) Extract(ctx context.Context, obj model.Obj, args model.ArchiveInnerArgs) (*model.Link, error) {
if !d.ForwardArchiveReq {
return nil, errs.NotSupport
}
var resp common.Resp[ArchiveMetaResp]
_, _, err := d.request("/fs/archive/meta", http.MethodPost, func(req *resty.Request) {
req.SetResult(&resp).SetBody(ArchiveMetaReq{
ArchivePass: args.Password,
Password: d.MetaPassword,
Path: obj.GetPath(),
Refresh: false,
})
})
if err != nil {
return nil, err
}
return &model.Link{
URL: fmt.Sprintf("%s?inner=%s&pass=%s&sign=%s",
resp.Data.RawURL,
utils.EncodePath(args.InnerPath, true),
url.QueryEscape(args.Password),
resp.Data.Sign),
}, nil
}
func (d *AListV3) ArchiveDecompress(ctx context.Context, srcObj, dstDir model.Obj, args model.ArchiveDecompressArgs) error {
if !d.ForwardArchiveReq {
return errs.NotImplement
}
dir, name := path.Split(srcObj.GetPath())
_, _, err := d.request("/fs/archive/decompress", http.MethodPost, func(req *resty.Request) {
req.SetBody(DecompressReq{
ArchivePass: args.Password,
CacheFull: args.CacheFull,
DstDir: dstDir.GetPath(),
InnerPath: args.InnerPath,
Name: []string{name},
PutIntoNewDir: args.PutIntoNewDir,
SrcDir: dir,
})
})
return err
}
func (d *AListV3) ResolveLinkCacheMode(_ string) driver.LinkCacheMode {
var mode driver.LinkCacheMode
if d.PassIPToUpsteam {
mode |= driver.LinkCacheIP
}
if d.PassUAToUpsteam {
mode |= driver.LinkCacheUA
}
return mode
}
var _ driver.Driver = (*AListV3)(nil)
+32
View File
@@ -0,0 +1,32 @@
package alist_v3
import (
"github.com/OpenListTeam/OpenList/v4/internal/driver"
"github.com/OpenListTeam/OpenList/v4/internal/op"
)
type Addition struct {
driver.RootPath
Address string `json:"url" required:"true"`
MetaPassword string `json:"meta_password"`
Username string `json:"username"`
Password string `json:"password"`
Token string `json:"token"`
PassIPToUpsteam bool `json:"pass_ip_to_upsteam" default:"true"`
PassUAToUpsteam bool `json:"pass_ua_to_upsteam" default:"true"`
ForwardArchiveReq bool `json:"forward_archive_requests" default:"true"`
}
var config = driver.Config{
Name: "AList V3",
LocalSort: true,
DefaultRoot: "/",
ProxyRangeOption: true,
LinkCacheMode: driver.LinkCacheAuto,
}
func init() {
op.RegisterDriver(func() driver.Driver {
return &AListV3{}
})
}
+170
View File
@@ -0,0 +1,170 @@
package alist_v3
import (
"time"
"github.com/OpenListTeam/OpenList/v4/internal/model"
"github.com/OpenListTeam/OpenList/v4/pkg/utils"
)
type ListReq struct {
model.PageReq
Path string `json:"path" form:"path"`
Password string `json:"password" form:"password"`
Refresh bool `json:"refresh"`
}
type ObjResp struct {
Name string `json:"name"`
Size int64 `json:"size"`
IsDir bool `json:"is_dir"`
Modified time.Time `json:"modified"`
Created time.Time `json:"created"`
Sign string `json:"sign"`
Thumb string `json:"thumb"`
Type int `json:"type"`
HashInfo string `json:"hashinfo"`
}
type FsListResp struct {
Content []ObjResp `json:"content"`
Total int64 `json:"total"`
Readme string `json:"readme"`
Write bool `json:"write"`
Provider string `json:"provider"`
}
type FsGetReq struct {
Path string `json:"path" form:"path"`
Password string `json:"password" form:"password"`
}
type FsGetResp struct {
ObjResp
RawURL string `json:"raw_url"`
Readme string `json:"readme"`
Provider string `json:"provider"`
Related []ObjResp `json:"related"`
}
type MkdirOrLinkReq struct {
Path string `json:"path" form:"path"`
}
type MoveCopyReq struct {
SrcDir string `json:"src_dir"`
DstDir string `json:"dst_dir"`
Names []string `json:"names"`
}
type RenameReq struct {
Path string `json:"path"`
Name string `json:"name"`
}
type RemoveReq struct {
Dir string `json:"dir"`
Names []string `json:"names"`
}
type LoginResp struct {
Token string `json:"token"`
}
type MeResp struct {
Id int `json:"id"`
Username string `json:"username"`
Password string `json:"password"`
BasePath string `json:"base_path"`
Role []int `json:"role"`
Disabled bool `json:"disabled"`
Permission int `json:"permission"`
SsoId string `json:"sso_id"`
Otp bool `json:"otp"`
}
type ArchiveMetaReq struct {
ArchivePass string `json:"archive_pass"`
Password string `json:"password"`
Path string `json:"path"`
Refresh bool `json:"refresh"`
}
type TreeResp struct {
ObjResp
Children []TreeResp `json:"children"`
hashCache *utils.HashInfo
}
func (t *TreeResp) GetSize() int64 {
return t.Size
}
func (t *TreeResp) GetName() string {
return t.Name
}
func (t *TreeResp) ModTime() time.Time {
return t.Modified
}
func (t *TreeResp) CreateTime() time.Time {
return t.Created
}
func (t *TreeResp) IsDir() bool {
return t.ObjResp.IsDir
}
func (t *TreeResp) GetHash() utils.HashInfo {
return utils.FromString(t.HashInfo)
}
func (t *TreeResp) GetID() string {
return ""
}
func (t *TreeResp) GetPath() string {
return ""
}
func (t *TreeResp) GetChildren() []model.ObjTree {
ret := make([]model.ObjTree, 0, len(t.Children))
for _, child := range t.Children {
ret = append(ret, &child)
}
return ret
}
func (t *TreeResp) Thumb() string {
return t.ObjResp.Thumb
}
type ArchiveMetaResp struct {
Comment string `json:"comment"`
Encrypted bool `json:"encrypted"`
Content []TreeResp `json:"content"`
RawURL string `json:"raw_url"`
Sign string `json:"sign"`
}
type ArchiveListReq struct {
model.PageReq
ArchiveMetaReq
InnerPath string `json:"inner_path"`
}
type ArchiveListResp struct {
Content []ObjResp `json:"content"`
Total int64 `json:"total"`
}
type DecompressReq struct {
ArchivePass string `json:"archive_pass"`
CacheFull bool `json:"cache_full"`
DstDir string `json:"dst_dir"`
InnerPath string `json:"inner_path"`
Name []string `json:"name"`
PutIntoNewDir bool `json:"put_into_new_dir"`
SrcDir string `json:"src_dir"`
}
+65
View File
@@ -0,0 +1,65 @@
package alist_v3
import (
"fmt"
"net/http"
"github.com/OpenListTeam/OpenList/v4/drivers/base"
"github.com/OpenListTeam/OpenList/v4/internal/op"
"github.com/OpenListTeam/OpenList/v4/pkg/utils"
"github.com/OpenListTeam/OpenList/v4/server/common"
"github.com/go-resty/resty/v2"
log "github.com/sirupsen/logrus"
)
func (d *AListV3) login() error {
if d.Username == "" {
return nil
}
var resp common.Resp[LoginResp]
_, _, err := d.request("/auth/login", http.MethodPost, func(req *resty.Request) {
req.SetResult(&resp).SetBody(base.Json{
"username": d.Username,
"password": d.Password,
})
})
if err != nil {
return err
}
d.Token = resp.Data.Token
op.MustSaveDriverStorage(d)
return nil
}
func (d *AListV3) request(api, method string, callback base.ReqCallback, retry ...bool) ([]byte, int, error) {
url := d.Address + "/api" + api
req := base.RestyClient.R()
req.SetHeader("Authorization", d.Token)
if callback != nil {
callback(req)
}
res, err := req.Execute(method, url)
if err != nil {
code := 0
if res != nil {
code = res.StatusCode()
}
return nil, code, err
}
log.Debugf("[openlist] response body: %s", res.String())
if res.StatusCode() >= 400 {
return nil, res.StatusCode(), fmt.Errorf("request failed, status: %s", res.Status())
}
code := utils.Json.Get(res.Body(), "code").ToInt()
if code != 200 {
if (code == 401 || code == 403) && !utils.IsBool(retry...) {
err = d.login()
if err != nil {
return nil, code, err
}
return d.request(api, method, callback, true)
}
return nil, code, fmt.Errorf("request failed,code: %d, message: %s", code, utils.Json.Get(res.Body(), "message").ToString())
}
return res.Body(), 200, nil
}
+1 -5
View File
@@ -77,7 +77,6 @@ func (d *AliyundriveOpen) GetRoot(ctx context.Context) (model.Obj, error) {
ID: d.RootFolderID,
Path: "/",
Name: "root",
Size: 0,
Modified: d.Modified,
IsFolder: true,
}, nil
@@ -299,10 +298,7 @@ func (d *AliyundriveOpen) GetDetails(ctx context.Context) (*model.StorageDetails
total := utils.Json.Get(res, "personal_space_info", "total_size").ToUint64()
used := utils.Json.Get(res, "personal_space_info", "used_size").ToUint64()
return &model.StorageDetails{
DiskUsage: model.DiskUsage{
TotalSpace: total,
FreeSpace: total - used,
},
DiskUsage: driver.DiskUsageFromUsedAndTotal(used, total),
}, nil
}
+2 -2
View File
@@ -242,11 +242,11 @@ func (d *AliyundriveOpen) upload(ctx context.Context, dstDir model.Obj, stream m
if err != nil {
return nil, err
}
rateLimitedRd := driver.NewLimitedUploadStream(ctx, rd)
err = retry.Do(func() error {
rd.Seek(0, io.SeekStart)
return d.uploadPart(ctx, rateLimitedRd, createResp.PartInfoList[i])
return d.uploadPart(ctx, driver.NewLimitedUploadStream(ctx, rd), createResp.PartInfoList[i])
},
retry.Context(ctx),
retry.Attempts(3),
retry.DelayType(retry.BackOffDelay),
retry.Delay(time.Second))
-1
View File
@@ -38,7 +38,6 @@ func (d *AliyundriveOpen) _refreshToken(ctx context.Context) (string, string, er
return "", "", err
}
_, err = base.RestyClient.R().
SetHeader("User-Agent", "Mozilla/5.0 (Macintosh; Apple macOS 15_5) AppleWebKit/537.36 (KHTML, like Gecko) Safari/537.36 Chrome/138.0.0.0 Openlist/425.6.30").
SetResult(&resp).
SetQueryParams(map[string]string{
"refresh_ui": d.RefreshToken,
+2 -1
View File
@@ -13,6 +13,7 @@ import (
_ "github.com/OpenListTeam/OpenList/v4/drivers/189_tv"
_ "github.com/OpenListTeam/OpenList/v4/drivers/189pc"
_ "github.com/OpenListTeam/OpenList/v4/drivers/alias"
_ "github.com/OpenListTeam/OpenList/v4/drivers/alist_v3"
_ "github.com/OpenListTeam/OpenList/v4/drivers/aliyundrive"
_ "github.com/OpenListTeam/OpenList/v4/drivers/aliyundrive_open"
_ "github.com/OpenListTeam/OpenList/v4/drivers/aliyundrive_share"
@@ -77,11 +78,11 @@ import (
_ "github.com/OpenListTeam/OpenList/v4/drivers/webdav"
_ "github.com/OpenListTeam/OpenList/v4/drivers/weiyun"
_ "github.com/OpenListTeam/OpenList/v4/drivers/wopan"
_ "github.com/OpenListTeam/OpenList/v4/drivers/wps"
_ "github.com/OpenListTeam/OpenList/v4/drivers/yandex_disk"
)
// All do nothing,just for import
// same as _ import
func All() {
}
+1 -2
View File
@@ -25,12 +25,11 @@ type AzureBlob struct {
Addition
client *azblob.Client
containerClient *container.Client
config driver.Config
}
// Config returns the driver configuration.
func (d *AzureBlob) Config() driver.Config {
return d.config
return config
}
// GetAddition returns additional settings specific to Azure Blob Storage.
+2 -8
View File
@@ -6,17 +6,13 @@ import (
)
type Addition struct {
driver.RootPath
Endpoint string `json:"endpoint" required:"true" default:"https://<accountname>.blob.core.windows.net/" help:"e.g. https://accountname.blob.core.windows.net/. The full endpoint URL for Azure Storage, including the unique storage account name (3 ~ 24 numbers and lowercase letters only)."`
AccessKey string `json:"access_key" required:"true" help:"The access key for Azure Storage, used for authentication. https://learn.microsoft.com/azure/storage/common/storage-account-keys-manage"`
ContainerName string `json:"container_name" required:"true" help:"The name of the container in Azure Storage (created in the Azure portal). https://learn.microsoft.com/azure/storage/blobs/blob-containers-portal"`
SignURLExpire int `json:"sign_url_expire" type:"number" default:"4" help:"The expiration time for SAS URLs, in hours."`
}
// implement GetRootId interface
func (r Addition) GetRootId() string {
return r.ContainerName
}
var config = driver.Config{
Name: "Azure Blob Storage",
LocalSort: true,
@@ -25,8 +21,6 @@ var config = driver.Config{
func init() {
op.RegisterDriver(func() driver.Driver {
return &AzureBlob{
config: config,
}
return &AzureBlob{}
})
}
+199 -74
View File
@@ -1,15 +1,19 @@
package baidu_netdisk
import (
"bytes"
"context"
"crypto/md5"
"encoding/hex"
"errors"
"io"
"mime/multipart"
"net/http"
"net/url"
"os"
stdpath "path"
"strconv"
"strings"
"time"
"github.com/OpenListTeam/OpenList/v4/drivers/base"
@@ -17,6 +21,7 @@ import (
"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/net"
"github.com/OpenListTeam/OpenList/v4/pkg/errgroup"
"github.com/OpenListTeam/OpenList/v4/pkg/utils"
"github.com/avast/retry-go"
@@ -31,6 +36,8 @@ type BaiduNetdisk struct {
vipType int // 会员类型,0普通用户(4G/4M)、1普通会员(10G/16M)、2超级会员(20G/32M)
}
var ErrUploadIDExpired = errors.New("uploadid expired")
func (d *BaiduNetdisk) Config() driver.Config {
return config
}
@@ -41,18 +48,20 @@ func (d *BaiduNetdisk) GetAddition() driver.Additional {
func (d *BaiduNetdisk) Init(ctx context.Context) error {
d.uploadThread, _ = strconv.Atoi(d.UploadThread)
if d.uploadThread < 1 || d.uploadThread > 32 {
d.uploadThread, d.UploadThread = 3, "3"
if d.uploadThread < 1 {
d.uploadThread, d.UploadThread = 1, "1"
} else if d.uploadThread > 32 {
d.uploadThread, d.UploadThread = 32, "32"
}
if _, err := url.Parse(d.UploadAPI); d.UploadAPI == "" || err != nil {
d.UploadAPI = "https://d.pcs.baidu.com"
d.UploadAPI = UPLOAD_FALLBACK_API
}
res, err := d.get("/xpan/nas", map[string]string{
"method": "uinfo",
}, nil)
log.Debugf("[baidu] get uinfo: %s", string(res))
log.Debugf("[baidu_netdisk] get uinfo: %s", string(res))
if err != nil {
return err
}
@@ -75,9 +84,10 @@ func (d *BaiduNetdisk) List(ctx context.Context, dir model.Obj, args model.ListA
}
func (d *BaiduNetdisk) Link(ctx context.Context, file model.Obj, args model.LinkArgs) (*model.Link, error) {
if d.DownloadAPI == "crack" {
switch d.DownloadAPI {
case "crack":
return d.linkCrack(file, args)
} else if d.DownloadAPI == "crack_video" {
case "crack_video":
return d.linkCrackVideo(file, args)
}
return d.linkOfficial(file, args)
@@ -179,6 +189,11 @@ func (d *BaiduNetdisk) PutRapid(ctx context.Context, dstDir model.Obj, stream mo
// **注意**: 截至 2024/04/20 百度云盘 api 接口返回的时间永远是当前时间,而不是文件时间。
// 而实际上云盘存储的时间是文件时间,所以此处需要覆盖时间,保证缓存与云盘的数据一致
func (d *BaiduNetdisk) Put(ctx context.Context, dstDir model.Obj, stream model.FileStreamer, up driver.UpdateProgress) (model.Obj, error) {
// 百度网盘不允许上传空文件
if stream.GetSize() < 1 {
return nil, ErrBaiduEmptyFilesNotAllowed
}
// rapid upload
if newObj, err := d.PutRapid(ctx, dstDir, stream); err == nil {
return newObj, nil
@@ -189,7 +204,7 @@ func (d *BaiduNetdisk) Put(ctx context.Context, dstDir model.Obj, stream model.F
tmpF *os.File
err error
)
if _, ok := cache.(io.ReaderAt); !ok {
if cache == nil {
tmpF, err = os.CreateTemp(conf.Conf.TempDir, "file-*")
if err != nil {
return nil, err
@@ -214,7 +229,6 @@ func (d *BaiduNetdisk) Put(ctx context.Context, dstDir model.Obj, stream model.F
// cal md5 for first 256k data
const SliceSize int64 = 256 * utils.KB
// cal md5
blockList := make([]string, 0, count)
byteSize := sliceSize
fileMd5H := md5.New()
@@ -244,7 +258,7 @@ func (d *BaiduNetdisk) Put(ctx context.Context, dstDir model.Obj, stream model.F
}
if tmpF != nil {
if written != streamSize {
return nil, errs.NewErr(err, "CreateTempFile failed, incoming stream actual size= %d, expect = %d ", written, streamSize)
return nil, errs.NewErr(err, "CreateTempFile failed, size mismatch: %d != %d ", written, streamSize)
}
_, err = tmpF.Seek(0, io.SeekStart)
if err != nil {
@@ -258,31 +272,14 @@ func (d *BaiduNetdisk) Put(ctx context.Context, dstDir model.Obj, stream model.F
mtime := stream.ModTime().Unix()
ctime := stream.CreateTime().Unix()
// step.1 预上传
// 尝试获取之前的进度
// step.1 尝试读取已保存进度
precreateResp, ok := base.GetUploadProgress[*PrecreateResp](d, d.AccessToken, contentMd5)
if !ok {
params := map[string]string{
"method": "precreate",
}
form := map[string]string{
"path": path,
"size": strconv.FormatInt(streamSize, 10),
"isdir": "0",
"autoinit": "1",
"rtype": "3",
"block_list": blockListStr,
"content-md5": contentMd5,
"slice-md5": sliceMd5,
}
joinTime(form, ctime, mtime)
log.Debugf("[baidu_netdisk] precreate data: %s", form)
_, err = d.postForm("/xpan/file", params, form, &precreateResp)
// 没有进度,走预上传
precreateResp, err = d.precreate(ctx, path, streamSize, blockListStr, contentMd5, sliceMd5, ctime, mtime)
if err != nil {
return nil, err
}
log.Debugf("%+v", precreateResp)
if precreateResp.ReturnType == 2 {
// rapid upload, since got md5 match from baidu server
// 修复时间,具体原因见 Put 方法注释的 **注意**
@@ -291,48 +288,95 @@ func (d *BaiduNetdisk) Put(ctx context.Context, dstDir model.Obj, stream model.F
return fileToObj(precreateResp.File), nil
}
}
// step.2 上传分片
threadG, upCtx := errgroup.NewGroupWithContext(ctx, d.uploadThread,
retry.Attempts(1),
retry.Delay(time.Second),
retry.DelayType(retry.BackOffDelay))
for i, partseq := range precreateResp.BlockList {
if utils.IsCanceled(upCtx) {
break
ensureUploadURL := func() {
if precreateResp.UploadURL != "" {
return
}
i, partseq, offset, byteSize := i, partseq, int64(partseq)*sliceSize, sliceSize
if partseq+1 == count {
byteSize = lastBlockSize
}
threadG.Go(func(ctx context.Context) error {
params := map[string]string{
"method": "upload",
"access_token": d.AccessToken,
"type": "tmpfile",
"path": path,
"uploadid": precreateResp.Uploadid,
"partseq": strconv.Itoa(partseq),
}
err := d.uploadSlice(ctx, params, stream.GetName(),
driver.NewLimitedUploadStream(ctx, io.NewSectionReader(cache, offset, byteSize)))
if err != nil {
return err
}
up(float64(threadG.Success()) * 100 / float64(len(precreateResp.BlockList)))
precreateResp.BlockList[i] = -1
return nil
})
precreateResp.UploadURL = d.getUploadUrl(path, precreateResp.Uploadid)
}
if err = threadG.Wait(); err != nil {
// 如果属于用户主动取消,则保存上传进度
// step.2 上传分片
uploadLoop:
for range 2 {
// 获取上传域名
ensureUploadURL()
// 并发上传
threadG, upCtx := errgroup.NewGroupWithContext(ctx, d.uploadThread,
retry.Attempts(UPLOAD_RETRY_COUNT),
retry.Delay(UPLOAD_RETRY_WAIT_TIME),
retry.MaxDelay(UPLOAD_RETRY_MAX_WAIT_TIME),
retry.DelayType(retry.BackOffDelay),
retry.RetryIf(func(err error) bool {
return !errors.Is(err, ErrUploadIDExpired)
}),
retry.LastErrorOnly(true))
totalParts := len(precreateResp.BlockList)
for i, partseq := range precreateResp.BlockList {
if utils.IsCanceled(upCtx) {
break
}
if partseq < 0 {
continue
}
i, partseq := i, partseq
offset, size := int64(partseq)*sliceSize, sliceSize
if partseq+1 == count {
size = lastBlockSize
}
threadG.Go(func(ctx context.Context) error {
params := map[string]string{
"method": "upload",
"access_token": d.AccessToken,
"type": "tmpfile",
"path": path,
"uploadid": precreateResp.Uploadid,
"partseq": strconv.Itoa(partseq),
}
section := io.NewSectionReader(cache, offset, size)
err := d.uploadSlice(ctx, precreateResp.UploadURL, params, stream.GetName(), section)
if err != nil {
return err
}
precreateResp.BlockList[i] = -1
progress := float64(threadG.Success()+1) * 100 / float64(totalParts+1)
up(progress)
return nil
})
}
err = threadG.Wait()
if err == nil {
break uploadLoop
}
// 保存进度(所有错误都会保存)
precreateResp.BlockList = utils.SliceFilter(precreateResp.BlockList, func(s int) bool { return s >= 0 })
base.SaveUploadProgress(d, precreateResp, d.AccessToken, contentMd5)
if errors.Is(err, context.Canceled) {
precreateResp.BlockList = utils.SliceFilter(precreateResp.BlockList, func(s int) bool { return s >= 0 })
return nil, err
}
if errors.Is(err, ErrUploadIDExpired) {
log.Warn("[baidu_netdisk] uploadid expired, will restart from scratch")
// 重新 precreate(所有分片都要重传)
newPre, err2 := d.precreate(ctx, path, streamSize, blockListStr, "", "", ctime, mtime)
if err2 != nil {
return nil, err2
}
if newPre.ReturnType == 2 {
return fileToObj(newPre.File), nil
}
precreateResp = newPre
precreateResp.UploadURL = ""
// 覆盖掉旧的进度
base.SaveUploadProgress(d, precreateResp, d.AccessToken, contentMd5)
continue uploadLoop
}
return nil, err
}
defer up(100)
// step.3 创建文件
var newFile File
@@ -343,23 +387,104 @@ func (d *BaiduNetdisk) Put(ctx context.Context, dstDir model.Obj, stream model.F
// 修复时间,具体原因见 Put 方法注释的 **注意**
newFile.Ctime = ctime
newFile.Mtime = mtime
// 上传成功清理进度
base.SaveUploadProgress(d, nil, d.AccessToken, contentMd5)
return fileToObj(newFile), nil
}
func (d *BaiduNetdisk) uploadSlice(ctx context.Context, params map[string]string, fileName string, file io.Reader) error {
res, err := base.RestyClient.R().
SetContext(ctx).
SetQueryParams(params).
SetFileReader("file", fileName, file).
Post(d.UploadAPI + "/rest/2.0/pcs/superfile2")
// precreate 执行预上传操作,支持首次上传和 uploadid 过期重试
func (d *BaiduNetdisk) precreate(ctx context.Context, path string, streamSize int64, blockListStr, contentMd5, sliceMd5 string, ctime, mtime int64) (*PrecreateResp, error) {
params := map[string]string{"method": "precreate"}
form := map[string]string{
"path": path,
"size": strconv.FormatInt(streamSize, 10),
"isdir": "0",
"autoinit": "1",
"rtype": "3",
"block_list": blockListStr,
}
// 只有在首次上传时才包含 content-md5 和 slice-md5
if contentMd5 != "" && sliceMd5 != "" {
form["content-md5"] = contentMd5
form["slice-md5"] = sliceMd5
}
joinTime(form, ctime, mtime)
var precreateResp PrecreateResp
_, err := d.postForm("/xpan/file", params, form, &precreateResp)
if err != nil {
return nil, err
}
// 修复时间,具体原因见 Put 方法注释的 **注意**
if precreateResp.ReturnType == 2 {
precreateResp.File.Ctime = ctime
precreateResp.File.Mtime = mtime
}
return &precreateResp, nil
}
func (d *BaiduNetdisk) uploadSlice(ctx context.Context, uploadUrl string, params map[string]string, fileName string, file *io.SectionReader) error {
b := bytes.NewBuffer(make([]byte, 0, bytes.MinRead))
mw := multipart.NewWriter(b)
_, err := mw.CreateFormFile("file", fileName)
if err != nil {
return err
}
log.Debugln(res.RawResponse.Status + res.String())
errCode := utils.Json.Get(res.Body(), "error_code").ToInt()
errNo := utils.Json.Get(res.Body(), "errno").ToInt()
headSize := b.Len()
err = mw.Close()
if err != nil {
return err
}
head := bytes.NewReader(b.Bytes()[:headSize])
tail := bytes.NewReader(b.Bytes()[headSize:])
rateLimitedRd := driver.NewLimitedUploadStream(ctx, io.MultiReader(head, file, tail))
req, err := http.NewRequestWithContext(ctx, http.MethodPost, uploadUrl+"/rest/2.0/pcs/superfile2", rateLimitedRd)
if err != nil {
return err
}
query := req.URL.Query()
for k, v := range params {
query.Set(k, v)
}
req.URL.RawQuery = query.Encode()
req.Header.Set("Content-Type", mw.FormDataContentType())
req.ContentLength = int64(b.Len()) + file.Size()
client := net.NewHttpClient()
if d.UploadSliceTimeout > 0 {
client.Timeout = time.Second * time.Duration(d.UploadSliceTimeout)
} else {
client.Timeout = DEFAULT_UPLOAD_SLICE_TIMEOUT
}
resp, err := client.Do(req)
if err != nil {
return err
}
defer resp.Body.Close()
b.Reset()
_, err = b.ReadFrom(resp.Body)
if err != nil {
return err
}
body := b.Bytes()
respStr := string(body)
log.Debugln(respStr)
lower := strings.ToLower(respStr)
// 合并 uploadid 过期检测逻辑
if strings.Contains(lower, "uploadid") &&
(strings.Contains(lower, "invalid") || strings.Contains(lower, "expired") || strings.Contains(lower, "not found")) {
return ErrUploadIDExpired
}
errCode := utils.Json.Get(body, "error_code").ToInt()
errNo := utils.Json.Get(body, "errno").ToInt()
if errCode != 0 || errNo != 0 {
return errs.NewErr(errs.StreamIncomplete, "error in uploading to baidu, will retry. response=%s", res.String())
return errs.NewErr(errs.StreamIncomplete, "error uploading to baidu, response=%s", respStr)
}
return nil
}
+12
View File
@@ -3,6 +3,7 @@ package baidu_netdisk
import (
"github.com/OpenListTeam/OpenList/v4/internal/driver"
"github.com/OpenListTeam/OpenList/v4/internal/op"
"time"
)
type Addition struct {
@@ -18,12 +19,23 @@ type Addition struct {
AccessToken string
RefreshToken string `json:"refresh_token" required:"true"`
UploadThread string `json:"upload_thread" default:"3" help:"1<=thread<=32"`
UploadSliceTimeout int `json:"upload_timeout" type:"number" default:"60" help:"per-slice upload timeout in seconds"`
UploadAPI string `json:"upload_api" default:"https://d.pcs.baidu.com"`
UseDynamicUploadAPI bool `json:"use_dynamic_upload_api" default:"true" help:"dynamically get upload api domain, when enabled, the 'Upload API' setting will be used as a fallback if failed to get"`
CustomUploadPartSize int64 `json:"custom_upload_part_size" type:"number" default:"0" help:"0 for auto"`
LowBandwithUploadMode bool `json:"low_bandwith_upload_mode" default:"false"`
OnlyListVideoFile bool `json:"only_list_video_file" default:"false"`
}
const (
UPLOAD_FALLBACK_API = "https://d.pcs.baidu.com" // 备用上传地址
UPLOAD_URL_EXPIRE_TIME = time.Minute * 60 // 上传地址有效期(分钟)
DEFAULT_UPLOAD_SLICE_TIMEOUT = time.Second * 60 // 上传分片请求默认超时时间
UPLOAD_RETRY_COUNT = 3
UPLOAD_RETRY_WAIT_TIME = time.Second * 1
UPLOAD_RETRY_MAX_WAIT_TIME = time.Second * 5
)
var config = driver.Config{
Name: "BaiduNetdisk",
DefaultRoot: "/",
+32 -4
View File
@@ -1,12 +1,16 @@
package baidu_netdisk
import (
"errors"
"path"
"strconv"
"time"
"github.com/OpenListTeam/OpenList/v4/internal/model"
"github.com/OpenListTeam/OpenList/v4/pkg/utils"
)
var (
ErrBaiduEmptyFilesNotAllowed = errors.New("empty files are not allowed by baidu netdisk")
)
type TokenErrResp struct {
@@ -71,9 +75,7 @@ func fileToObj(f File) *model.ObjThumb {
Modified: time.Unix(f.ServerMtime, 0),
Ctime: time.Unix(f.ServerCtime, 0),
IsFolder: f.Isdir == 1,
// 直接获取的MD5是错误的
HashInfo: utils.NewHashInfo(utils.MD5, DecryptMd5(f.Md5)),
// 百度API返回的MD5不可信,不使用HashInfo
},
Thumbnail: model.Thumbnail{Thumbnail: f.Thumbs.Url3},
}
@@ -188,6 +190,32 @@ type PrecreateResp struct {
// return_type=2
File File `json:"info"`
UploadURL string `json:"-"` // 保存断点续传对应的上传域名
}
type UploadServerResp struct {
BakServer []any `json:"bak_server"`
BakServers []struct {
Server string `json:"server"`
} `json:"bak_servers"`
ClientIP string `json:"client_ip"`
ErrorCode int `json:"error_code"`
ErrorMsg string `json:"error_msg"`
Expire int `json:"expire"`
Host string `json:"host"`
Newno string `json:"newno"`
QuicServer []any `json:"quic_server"`
QuicServers []struct {
Server string `json:"server"`
} `json:"quic_servers"`
RequestID int64 `json:"request_id"`
Server []any `json:"server"`
ServerTime int `json:"server_time"`
Servers []struct {
Server string `json:"server"`
} `json:"servers"`
Sl int `json:"sl"`
}
type QuotaResp struct {
+55 -10
View File
@@ -42,7 +42,6 @@ func (d *BaiduNetdisk) _refreshToken() error {
ErrorMessage string `json:"text"`
}
_, err := base.RestyClient.R().
SetHeader("User-Agent", "Mozilla/5.0 (Macintosh; Apple macOS 15_5) AppleWebKit/537.36 (KHTML, like Gecko) Safari/537.36 Chrome/138.0.0.0 Openlist/425.6.30").
SetResult(&resp).
SetQueryParams(map[string]string{
"refresh_ui": d.RefreshToken,
@@ -115,14 +114,14 @@ func (d *BaiduNetdisk) request(furl string, method string, callback base.ReqCall
errno := utils.Json.Get(res.Body(), "errno").ToInt()
if errno != 0 {
if utils.SliceContains([]int{111, -6}, errno) {
log.Info("refreshing baidu_netdisk token.")
log.Info("[baidu_netdisk] refreshing baidu_netdisk token.")
err2 := d.refreshToken()
if err2 != nil {
return retry.Unrecoverable(err2)
}
}
if 31023 == errno && d.DownloadAPI == "crack_video" {
if errno == 31023 && d.DownloadAPI == "crack_video" {
result = res.Body()
return nil
}
@@ -247,7 +246,7 @@ func (d *BaiduNetdisk) linkCrack(file model.Obj, _ model.LinkArgs) (*model.Link,
func (d *BaiduNetdisk) linkCrackVideo(file model.Obj, _ model.LinkArgs) (*model.Link, error) {
param := map[string]string{
"type": "VideoURL",
"path": fmt.Sprintf("%s", file.GetPath()),
"path": file.GetPath(),
"fs_id": file.GetID(),
"devuid": "0%1",
"clienttype": "1",
@@ -326,10 +325,10 @@ func (d *BaiduNetdisk) getSliceSize(filesize int64) int64 {
// 非会员固定为 4MB
if d.vipType == 0 {
if d.CustomUploadPartSize != 0 {
log.Warnf("CustomUploadPartSize is not supported for non-vip user, use DefaultSliceSize")
log.Warnf("[baidu_netdisk] CustomUploadPartSize is not supported for non-vip user, use DefaultSliceSize")
}
if filesize > MaxSliceNum*DefaultSliceSize {
log.Warnf("File size(%d) is too large, may cause upload failure", filesize)
log.Warnf("[baidu_netdisk] File size(%d) is too large, may cause upload failure", filesize)
}
return DefaultSliceSize
@@ -337,17 +336,17 @@ func (d *BaiduNetdisk) getSliceSize(filesize int64) int64 {
if d.CustomUploadPartSize != 0 {
if d.CustomUploadPartSize < DefaultSliceSize {
log.Warnf("CustomUploadPartSize(%d) is less than DefaultSliceSize(%d), use DefaultSliceSize", d.CustomUploadPartSize, DefaultSliceSize)
log.Warnf("[baidu_netdisk] CustomUploadPartSize(%d) is less than DefaultSliceSize(%d), use DefaultSliceSize", d.CustomUploadPartSize, DefaultSliceSize)
return DefaultSliceSize
}
if d.vipType == 1 && d.CustomUploadPartSize > VipSliceSize {
log.Warnf("CustomUploadPartSize(%d) is greater than VipSliceSize(%d), use VipSliceSize", d.CustomUploadPartSize, VipSliceSize)
log.Warnf("[baidu_netdisk] CustomUploadPartSize(%d) is greater than VipSliceSize(%d), use VipSliceSize", d.CustomUploadPartSize, VipSliceSize)
return VipSliceSize
}
if d.vipType == 2 && d.CustomUploadPartSize > SVipSliceSize {
log.Warnf("CustomUploadPartSize(%d) is greater than SVipSliceSize(%d), use SVipSliceSize", d.CustomUploadPartSize, SVipSliceSize)
log.Warnf("[baidu_netdisk] CustomUploadPartSize(%d) is greater than SVipSliceSize(%d), use SVipSliceSize", d.CustomUploadPartSize, SVipSliceSize)
return SVipSliceSize
}
@@ -377,7 +376,7 @@ func (d *BaiduNetdisk) getSliceSize(filesize int64) int64 {
}
if filesize > MaxSliceNum*maxSliceSize {
log.Warnf("File size(%d) is too large, may cause upload failure", filesize)
log.Warnf("[baidu_netdisk] File size(%d) is too large, may cause upload failure", filesize)
}
return maxSliceSize
@@ -394,6 +393,52 @@ func (d *BaiduNetdisk) quota(ctx context.Context) (model.DiskUsage, error) {
return driver.DiskUsageFromUsedAndTotal(resp.Used, resp.Total), nil
}
// getUploadUrl 从开放平台获取上传域名/地址,并发请求会被合并,结果会在 uploadid 生命周期内复用。
// 如果获取失败,则返回 Upload API设置项。
func (d *BaiduNetdisk) getUploadUrl(path, uploadId string) string {
if !d.UseDynamicUploadAPI || uploadId == "" {
return d.UploadAPI
}
uploadUrl, err := d.requestForUploadUrl(path, uploadId)
if err != nil {
return d.UploadAPI
}
return uploadUrl
}
// requestForUploadUrl 请求获取上传地址。
// 实测此接口不需要认证,传method和upload_version就行,不过还是按文档规范调用。
// https://pan.baidu.com/union/doc/Mlvw5hfnr
func (d *BaiduNetdisk) requestForUploadUrl(path, uploadId string) (string, error) {
params := map[string]string{
"method": "locateupload",
"appid": "250528",
"path": path,
"uploadid": uploadId,
"upload_version": "2.0",
}
apiUrl := "https://d.pcs.baidu.com/rest/2.0/pcs/file"
var resp UploadServerResp
_, err := d.request(apiUrl, http.MethodGet, func(req *resty.Request) {
req.SetQueryParams(params)
}, &resp)
if err != nil {
return "", err
}
// 应该是https开头的一个地址
var uploadUrl string
if len(resp.Servers) > 0 {
uploadUrl = resp.Servers[0].Server
} else if len(resp.BakServers) > 0 {
uploadUrl = resp.BakServers[0].Server
}
if uploadUrl == "" {
return "", errors.New("upload URL is empty")
}
return uploadUrl, nil
}
// func encodeURIComponent(str string) string {
// r := url.QueryEscape(str)
// r = strings.ReplaceAll(r, "+", "%20")
+2 -1
View File
@@ -371,7 +371,7 @@ func (d *BaiduPhoto) Put(ctx context.Context, dstDir model.Obj, stream model.Fil
if err != nil {
return err
}
up(float64(threadG.Success()) * 100 / float64(len(precreateResp.BlockList)))
up(float64(threadG.Success()+1) * 100 / float64(len(precreateResp.BlockList)+1))
precreateResp.BlockList[i] = -1
return nil
})
@@ -383,6 +383,7 @@ func (d *BaiduPhoto) Put(ctx context.Context, dstDir model.Obj, stream model.Fil
}
return nil, err
}
defer up(100)
fallthrough
case 2: //step.4 创建文件
params["uploadid"] = precreateResp.UploadID
+1 -1
View File
@@ -20,7 +20,7 @@ type Addition struct {
var config = driver.Config{
Name: "BaiduPhoto",
LocalSort: true,
LinkCacheType: 2,
LinkCacheMode: driver.LinkCacheUA,
}
func init() {
+7 -1
View File
@@ -15,9 +15,12 @@ var (
RestyClient *resty.Client
HttpClient *http.Client
)
var UserAgent = "Mozilla/5.0 (Macintosh; Apple macOS 15_5) AppleWebKit/537.36 (KHTML, like Gecko) Safari/537.36 Chrome/138.0.0.0"
var DefaultTimeout = time.Second * 30
const UserAgent = "Mozilla/5.0 (Macintosh; Apple macOS 26_1_0) AppleWebKit/537.36 (KHTML, like Gecko) Safari/537.36 Chrome/142.0.0.0 OpenList/425.6.30"
const UserAgentNT = "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Safari/537.36 Chrome/142.0.0.0 OpenList/425.6.30"
func InitClient() {
NoRedirectClient = resty.New().SetRedirectPolicy(
resty.RedirectPolicyFunc(func(req *http.Request, via []*http.Request) error {
@@ -25,6 +28,7 @@ func InitClient() {
}),
).SetTLSClientConfig(&tls.Config{InsecureSkipVerify: conf.Conf.TlsInsecureSkipVerify})
NoRedirectClient.SetHeader("user-agent", UserAgent)
net.SetRestyProxyIfConfigured(NoRedirectClient)
RestyClient = NewRestyClient()
HttpClient = net.NewHttpClient()
@@ -37,5 +41,7 @@ func NewRestyClient() *resty.Client {
SetRetryResetReaders(true).
SetTimeout(DefaultTimeout).
SetTLSClientConfig(&tls.Config{InsecureSkipVerify: conf.Conf.TlsInsecureSkipVerify})
net.SetRestyProxyIfConfigured(client)
return client
}
+12 -11
View File
@@ -226,16 +226,13 @@ func (d *ChaoXing) Put(ctx context.Context, dstDir model.Obj, file model.FileStr
if resp.Result != 1 {
return errors.New("get upload data error")
}
body := &bytes.Buffer{}
body := bytes.NewBuffer(make([]byte, 0, bytes.MinRead))
writer := multipart.NewWriter(body)
filePart, err := writer.CreateFormFile("file", file.GetName())
if err != nil {
return err
}
_, err = utils.CopyWithBuffer(filePart, file)
_, err = writer.CreateFormFile("file", file.GetName())
if err != nil {
return err
}
headSize := body.Len()
err = writer.WriteField("_token", resp.Msg.Token)
if err != nil {
return err
@@ -249,30 +246,34 @@ func (d *ChaoXing) Put(ctx context.Context, dstDir model.Obj, file model.FileStr
if err != nil {
return err
}
head := bytes.NewReader(body.Bytes()[:headSize])
tail := bytes.NewReader(body.Bytes()[headSize:])
r := driver.NewLimitedUploadStream(ctx, &driver.ReaderUpdatingProgress{
Reader: &driver.SimpleReaderWithSize{
Reader: body,
Size: int64(body.Len()),
Reader: io.MultiReader(head, file, tail),
Size: int64(body.Len()) + file.GetSize(),
},
UpdateProgress: up,
})
req, err := http.NewRequestWithContext(ctx, http.MethodPost, "https://pan-yz.chaoxing.com/upload", r)
if err != nil {
return err
}
req.Header.Set("Content-Type", writer.FormDataContentType())
req.Header.Set("Content-Length", strconv.Itoa(body.Len()))
req.ContentLength = int64(body.Len()) + file.GetSize()
resps, err := http.DefaultClient.Do(req)
if err != nil {
return err
}
defer resps.Body.Close()
bodys, err := io.ReadAll(resps.Body)
body.Reset()
_, err = body.ReadFrom(resps.Body)
if err != nil {
return err
}
var fileRsp UploadFileDataRsp
err = json.Unmarshal(bodys, &fileRsp)
err = json.Unmarshal(body.Bytes(), &fileRsp)
if err != nil {
return err
}
+15 -20
View File
@@ -10,6 +10,7 @@ import (
"strconv"
"strings"
"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/fs"
@@ -52,14 +53,11 @@ func (d *Chunk) Drop(ctx context.Context) error {
return nil
}
func (Addition) GetRootPath() string {
return "/"
}
func (d *Chunk) Get(ctx context.Context, path string) (model.Obj, error) {
if utils.PathEqual(path, "/") {
return &model.Object{
Name: "Root",
IsFolder: true,
Path: "/",
}, nil
}
remoteStorage, remoteActualPath, err := op.GetStorageAndActualPath(d.RemotePath)
if err != nil {
return nil, err
@@ -272,17 +270,13 @@ func (d *Chunk) Link(ctx context.Context, file model.Obj, args model.LinkArgs) (
// 检查0号块不等于-1 以支持空文件
// 如果块数量大于1 最后一块不可能为0
// 只检查中间块是否有0
for i, l := 0, len(chunkFile.chunkSizes)-2; ; i++ {
if i == 0 {
if chunkFile.chunkSizes[i] == -1 {
return nil, fmt.Errorf("chunk part[%d] are missing", i)
}
} else if chunkFile.chunkSizes[i] == 0 {
if chunkFile.chunkSizes[0] == -1 {
return nil, fmt.Errorf("chunk part[%d] are missing", 0)
}
for i, l := 1, len(chunkFile.chunkSizes)-1; i < l; i++ {
if chunkFile.chunkSizes[i] == 0 {
return nil, fmt.Errorf("chunk part[%d] are missing", i)
}
if i >= l {
break
}
}
fileSize := chunkFile.GetSize()
mergedRrf := func(ctx context.Context, httpRange http_range.Range) (io.ReadCloser, error) {
@@ -429,9 +423,10 @@ func (d *Chunk) Put(ctx context.Context, dstDir model.Obj, file model.FileStream
UpdateProgress: up,
}
dst := stdpath.Join(remoteActualPath, dstDir.GetPath(), d.ChunkPrefix+file.GetName())
skipHookCtx := context.WithValue(ctx, conf.SkipHookKey, struct{}{})
if d.StoreHash {
for ht, value := range file.GetHash().All() {
_ = op.Put(ctx, remoteStorage, dst, &stream.FileStream{
_ = op.Put(skipHookCtx, remoteStorage, dst, &stream.FileStream{
Obj: &model.Object{
Name: fmt.Sprintf("hash_%s_%s%s", ht.Name, value, d.CustomExt),
Size: 1,
@@ -439,7 +434,7 @@ func (d *Chunk) Put(ctx context.Context, dstDir model.Obj, file model.FileStream
},
Mimetype: "application/octet-stream",
Reader: bytes.NewReader([]byte{0}), // 兼容不支持空文件的驱动
}, nil, true)
}, nil)
}
}
fullPartCount := int(file.GetSize() / d.PartSize)
@@ -450,7 +445,7 @@ func (d *Chunk) Put(ctx context.Context, dstDir model.Obj, file model.FileStream
}
partIndex := 0
for partIndex < fullPartCount {
err = op.Put(ctx, remoteStorage, dst, &stream.FileStream{
err = op.Put(skipHookCtx, remoteStorage, dst, &stream.FileStream{
Obj: &model.Object{
Name: d.getPartName(partIndex),
Size: d.PartSize,
@@ -458,7 +453,7 @@ func (d *Chunk) Put(ctx context.Context, dstDir model.Obj, file model.FileStream
},
Mimetype: file.GetMimetype(),
Reader: io.LimitReader(upReader, d.PartSize),
}, nil, true)
}, nil)
if err != nil {
_ = op.Remove(ctx, remoteStorage, dst)
return err
+7 -2
View File
@@ -291,6 +291,7 @@ func (d *Cloudreve) upRemote(ctx context.Context, stream model.FileStreamer, u U
}
return nil
},
retry.Context(ctx),
retry.Attempts(3),
retry.DelayType(retry.BackOffDelay),
retry.Delay(time.Second),
@@ -351,7 +352,9 @@ func (d *Cloudreve) upOneDrive(ctx context.Context, stream model.FileStreamer, u
default:
return nil
}
}, retry.Attempts(3),
},
retry.Context(ctx),
retry.Attempts(3),
retry.DelayType(retry.BackOffDelay),
retry.Delay(time.Second),
)
@@ -414,7 +417,9 @@ func (d *Cloudreve) upS3(ctx context.Context, stream model.FileStreamer, u Uploa
etags = append(etags, etag)
return nil
}
}, retry.Attempts(3),
},
retry.Context(ctx),
retry.Attempts(3),
retry.DelayType(retry.BackOffDelay),
retry.Delay(time.Second),
)
+3 -1
View File
@@ -299,7 +299,9 @@ func (d *CloudreveV4) Put(ctx context.Context, dstDir model.Obj, file model.File
case "onedrive":
err = d.upOneDrive(ctx, file, u, up)
case "s3":
err = d.upS3(ctx, file, u, up)
err = d.upS3(ctx, file, u, up, "s3")
case "ks3":
err = d.upS3(ctx, file, u, up, "ks3")
default:
return errs.NotImplement
}
+17 -7
View File
@@ -447,7 +447,9 @@ func (d *CloudreveV4) upRemote(ctx context.Context, file model.FileStreamer, u F
return errors.New(up.Msg)
}
return nil
}, retry.Attempts(3),
},
retry.Context(ctx),
retry.Attempts(3),
retry.DelayType(retry.BackOffDelay),
retry.Delay(time.Second),
)
@@ -508,7 +510,9 @@ func (d *CloudreveV4) upOneDrive(ctx context.Context, file model.FileStreamer, u
default:
return nil
}
}, retry.Attempts(3),
},
retry.Context(ctx),
retry.Attempts(3),
retry.DelayType(retry.BackOffDelay),
retry.Delay(time.Second),
)
@@ -525,7 +529,7 @@ func (d *CloudreveV4) upOneDrive(ctx context.Context, file model.FileStreamer, u
}, nil)
}
func (d *CloudreveV4) upS3(ctx context.Context, file model.FileStreamer, u FileUploadResp, up driver.UpdateProgress) error {
func (d *CloudreveV4) upS3(ctx context.Context, file model.FileStreamer, u FileUploadResp, up driver.UpdateProgress, s3Type string) error {
DEFAULT := int64(u.ChunkSize)
ss, err := stream.NewStreamSectionReader(file, int(DEFAULT), &up)
if err != nil {
@@ -556,6 +560,9 @@ func (d *CloudreveV4) upS3(ctx context.Context, file model.FileStreamer, u FileU
}
req.ContentLength = byteSize
req.Header.Set("User-Agent", d.getUA())
if s3Type == "ks3" {
req.Header.Set("Content-Type", "application/octet-stream")
}
res, err := base.HttpClient.Do(req)
if err != nil {
return err
@@ -572,6 +579,7 @@ func (d *CloudreveV4) upS3(ctx context.Context, file model.FileStreamer, u FileU
return nil
}
},
retry.Context(ctx),
retry.Attempts(3),
retry.DelayType(retry.BackOffDelay),
retry.Delay(time.Second),
@@ -604,7 +612,11 @@ func (d *CloudreveV4) upS3(ctx context.Context, file model.FileStreamer, u FileU
if err != nil {
return err
}
req.Header.Set("Content-Type", "application/xml")
if s3Type == "ks3" {
req.Header.Set("Content-Type", "application/octet-stream")
} else {
req.Header.Set("Content-Type", "application/xml")
}
req.Header.Set("User-Agent", d.getUA())
res, err := base.HttpClient.Do(req)
if err != nil {
@@ -617,7 +629,5 @@ func (d *CloudreveV4) upS3(ctx context.Context, file model.FileStreamer, u FileU
}
// 上传成功发送回调请求
return d.request(http.MethodGet, "/callback/s3/"+u.SessionID+"/"+u.CallbackSecret, func(req *resty.Request) {
req.SetBody("{}")
}, nil)
return d.request(http.MethodGet, "/callback/"+s3Type+"/"+u.SessionID+"/"+u.CallbackSecret, nil, nil)
}
+35 -33
View File
@@ -50,7 +50,8 @@ func (d *CnbReleases) Drop(ctx context.Context) error {
}
func (d *CnbReleases) List(ctx context.Context, dir model.Obj, args model.ListArgs) ([]model.Obj, error) {
if dir.GetPath() == "/" {
dirID := dir.GetID()
if dirID == "" {
// get all releases for root dir
var resp ReleaseList
@@ -75,36 +76,32 @@ func (d *CnbReleases) List(ctx context.Context, dir model.Obj, args model.ListAr
IsFolder: true,
}, nil
})
} else {
// get release info by release id
releaseID := dir.GetID()
if releaseID == "" {
return nil, errs.ObjectNotFound
}
var resp Release
err := d.Request(http.MethodGet, "/{repo}/-/releases/{release_id}", func(req *resty.Request) {
req.SetPathParam("repo", d.Repo)
req.SetPathParam("release_id", releaseID)
}, &resp)
if err != nil {
return nil, err
}
return utils.SliceConvert(resp.Assets, func(src ReleaseAsset) (model.Obj, error) {
return &Object{
Object: model.Object{
ID: src.ID,
Path: src.Path,
Name: src.Name,
Size: src.Size,
Ctime: src.CreatedAt,
Modified: src.UpdatedAt,
IsFolder: false,
},
ParentID: dir.GetID(),
}, nil
})
}
var resp Release
err := d.Request(http.MethodGet, "/{repo}/-/releases/{release_id}", func(req *resty.Request) {
req.SetPathParam("repo", d.Repo)
req.SetPathParam("release_id", dirID)
}, &resp)
if err != nil {
return nil, err
}
return utils.SliceConvert(resp.Assets, func(src ReleaseAsset) (model.Obj, error) {
return &Object{
Object: model.Object{
ID: src.ID,
Path: src.Path,
Name: src.Name,
Size: src.Size,
Ctime: src.CreatedAt,
Modified: src.UpdatedAt,
IsFolder: false,
},
ParentID: dirID,
}, nil
})
}
func (d *CnbReleases) Link(ctx context.Context, file model.Obj, args model.LinkArgs) (*model.Link, error) {
@@ -200,15 +197,20 @@ func (d *CnbReleases) Put(ctx context.Context, dstDir model.Obj, file model.File
if err != nil {
return err
}
head := bytes.NewReader(b.Bytes()[:headSize])
tail := bytes.NewReader(b.Bytes()[headSize:])
rateLimitedRd := driver.NewLimitedUploadStream(ctx, io.MultiReader(head, file, tail))
r := driver.NewLimitedUploadStream(ctx, &driver.ReaderUpdatingProgress{
Reader: &driver.SimpleReaderWithSize{
Reader: io.MultiReader(head, file, tail),
Size: int64(b.Len()) + file.GetSize(),
},
UpdateProgress: up,
})
// use net/http to upload file
ctxWithTimeout, cancel := context.WithTimeout(ctx, time.Duration(resp.ExpiresInSec+1)*time.Second)
defer cancel()
req, err := http.NewRequestWithContext(ctxWithTimeout, http.MethodPost, resp.UploadURL, rateLimitedRd)
req, err := http.NewRequestWithContext(ctxWithTimeout, http.MethodPost, resp.UploadURL, r)
if err != nil {
return err
}
+3 -4
View File
@@ -6,7 +6,7 @@ import (
)
type Addition struct {
driver.RootPath
driver.RootID
Repo string `json:"repo" type:"string" required:"true"`
Token string `json:"token" type:"string" required:"true"`
UseTagName bool `json:"use_tag_name" type:"bool" default:"false" help:"Use tag name instead of release name"`
@@ -14,9 +14,8 @@ type Addition struct {
}
var config = driver.Config{
Name: "CNB Releases",
LocalSort: true,
DefaultRoot: "/",
Name: "CNB Releases",
LocalSort: true,
}
func init() {
+128 -158
View File
@@ -3,6 +3,7 @@ package crypt
import (
"bytes"
"context"
"errors"
"fmt"
"io"
stdpath "path"
@@ -29,8 +30,7 @@ import (
type Crypt struct {
model.Storage
Addition
cipher *rcCrypt.Cipher
remoteStorage driver.Driver
cipher *rcCrypt.Cipher
}
const obfuscatedPrefix = "___Obfuscated___"
@@ -60,15 +60,7 @@ func (d *Crypt) Init(ctx context.Context) error {
}
d.FileNameEncoding = utils.GetNoneEmpty(d.FileNameEncoding, "base64")
d.EncryptedSuffix = utils.GetNoneEmpty(d.EncryptedSuffix, ".bin")
op.MustSaveDriverStorage(d)
// need remote storage exist
storage, err := fs.GetStorage(d.RemotePath, &fs.GetStoragesArgs{})
if err != nil {
return fmt.Errorf("can't find remote storage: %w", err)
}
d.remoteStorage = storage
d.RemotePath = utils.FixAndCleanPath(d.RemotePath)
p, _ := strings.CutPrefix(d.Password, obfuscatedPrefix)
p2, _ := strings.CutPrefix(d.Salt, obfuscatedPrefix)
@@ -108,150 +100,146 @@ func (d *Crypt) Drop(ctx context.Context) error {
}
func (d *Crypt) List(ctx context.Context, dir model.Obj, args model.ListArgs) ([]model.Obj, error) {
path := dir.GetPath()
// return d.list(ctx, d.RemotePath, path)
// remoteFull
objs, err := fs.List(ctx, d.getPathForRemote(path, true), &fs.ListArgs{NoLog: true, Refresh: args.Refresh})
remoteFullPath := dir.GetPath()
objs, err := fs.List(ctx, remoteFullPath, &fs.ListArgs{NoLog: true, Refresh: args.Refresh})
// the obj must implement the model.SetPath interface
// return objs, err
if err != nil {
return nil, err
}
var result []model.Obj
result := make([]model.Obj, 0, len(objs))
for _, obj := range objs {
if obj.IsDir() {
name, err := d.cipher.DecryptDirName(obj.GetName())
if err != nil {
// filter illegal files
continue
}
if !d.ShowHidden && strings.HasPrefix(name, ".") {
continue
}
objRes := model.Object{
Name: name,
Size: 0,
Modified: obj.ModTime(),
IsFolder: obj.IsDir(),
Ctime: obj.CreateTime(),
// discarding hash as it's encrypted
}
result = append(result, &objRes)
} else {
thumb, ok := model.GetThumb(obj)
size, err := d.cipher.DecryptedSize(obj.GetSize())
if err != nil {
// filter illegal files
continue
}
name, err := d.cipher.DecryptFileName(obj.GetName())
if err != nil {
// filter illegal files
continue
}
if !d.ShowHidden && strings.HasPrefix(name, ".") {
continue
}
objRes := model.Object{
Name: name,
Size: size,
Modified: obj.ModTime(),
IsFolder: obj.IsDir(),
Ctime: obj.CreateTime(),
// discarding hash as it's encrypted
}
if d.Thumbnail && thumb == "" {
thumbPath := stdpath.Join(args.ReqPath, ".thumbnails", name+".webp")
thumb = fmt.Sprintf("%s/d%s?sign=%s",
common.GetApiUrl(ctx),
utils.EncodePath(thumbPath, true),
sign.Sign(thumbPath))
}
if !ok && !d.Thumbnail {
result = append(result, &objRes)
} else {
objWithThumb := model.ObjThumb{
Object: objRes,
Thumbnail: model.Thumbnail{
Thumbnail: thumb,
},
size := obj.GetSize()
mask := model.GetObjMask(obj)
name := obj.GetName()
if mask&model.Virtual == 0 {
if obj.IsDir() {
name, err = d.cipher.DecryptDirName(model.UnwrapObjName(obj).GetName())
if err != nil {
// filter illegal files
continue
}
} else {
size, err = d.cipher.DecryptedSize(size)
if err != nil {
// filter illegal files
continue
}
name, err = d.cipher.DecryptFileName(model.UnwrapObjName(obj).GetName())
if err != nil {
// filter illegal files
continue
}
result = append(result, &objWithThumb)
}
}
if !d.ShowHidden && strings.HasPrefix(name, ".") {
continue
}
objRes := &model.Object{
Path: stdpath.Join(remoteFullPath, obj.GetName()),
Name: name,
Size: size,
Modified: obj.ModTime(),
IsFolder: obj.IsDir(),
Ctime: obj.CreateTime(),
Mask: mask &^ model.Temp,
// discarding hash as it's encrypted
}
if !d.Thumbnail || !strings.HasPrefix(args.ReqPath, "/") {
result = append(result, objRes)
continue
}
thumbPath := stdpath.Join(args.ReqPath, ".thumbnails", name+".webp")
thumb := fmt.Sprintf("%s/d%s?sign=%s",
common.GetApiUrl(ctx),
utils.EncodePath(thumbPath, true),
sign.Sign(thumbPath))
result = append(result, &model.ObjThumb{
Object: *objRes,
Thumbnail: model.Thumbnail{
Thumbnail: thumb,
},
})
}
return result, nil
}
func (a Addition) GetRootPath() string {
return a.RemotePath
}
func (d *Crypt) Get(ctx context.Context, path string) (model.Obj, error) {
if utils.PathEqual(path, "/") {
return &model.Object{
Name: "Root",
IsFolder: true,
Path: "/",
}, nil
}
remoteFullPath := ""
var remoteObj model.Obj
var err, err2 error
firstTryIsFolder, secondTry := guessPath(path)
remoteFullPath = d.getPathForRemote(path, firstTryIsFolder)
remoteObj, err = fs.Get(ctx, remoteFullPath, &fs.GetArgs{NoLog: true})
remoteFullPath := stdpath.Join(d.RemotePath, d.encryptPath(path, firstTryIsFolder))
remoteObj, err := fs.Get(ctx, remoteFullPath, &fs.GetArgs{NoLog: true})
if err != nil {
if errs.IsObjectNotFound(err) && secondTry {
if errors.Is(err, errs.StorageNotFound) {
remoteFullPath = stdpath.Join(d.RemotePath, path)
remoteObj, err = fs.Get(ctx, remoteFullPath, &fs.GetArgs{NoLog: true})
if err != nil {
// 可能是 虚拟路径+开启文件夹加密:返回NotSupport让op.Get去尝试op.List查找
return nil, errs.NotSupport
}
} else if secondTry && errs.IsObjectNotFound(err) {
// try the opposite
remoteFullPath = d.getPathForRemote(path, !firstTryIsFolder)
remoteObj, err2 = fs.Get(ctx, remoteFullPath, &fs.GetArgs{NoLog: true})
if err2 != nil {
return nil, err2
remoteFullPath = stdpath.Join(d.RemotePath, d.encryptPath(path, !firstTryIsFolder))
remoteObj, err = fs.Get(ctx, remoteFullPath, &fs.GetArgs{NoLog: true})
if err != nil {
return nil, err
}
} else {
return nil, err
}
}
var size int64 = 0
name := ""
if !remoteObj.IsDir() {
size, err = d.cipher.DecryptedSize(remoteObj.GetSize())
if err != nil {
log.Warnf("DecryptedSize failed for %s ,will use original size, err:%s", path, err)
size = remoteObj.GetSize()
}
name, err = d.cipher.DecryptFileName(remoteObj.GetName())
if err != nil {
log.Warnf("DecryptFileName failed for %s ,will use original name, err:%s", path, err)
name = remoteObj.GetName()
}
} else {
name, err = d.cipher.DecryptDirName(remoteObj.GetName())
if err != nil {
log.Warnf("DecryptDirName failed for %s ,will use original name, err:%s", path, err)
name = remoteObj.GetName()
size := remoteObj.GetSize()
name := remoteObj.GetName()
mask := model.GetObjMask(remoteObj) &^ model.Temp
if mask&model.Virtual == 0 {
if !remoteObj.IsDir() {
decryptedSize, err := d.cipher.DecryptedSize(size)
if err != nil {
log.Warnf("DecryptedSize failed for %s ,will use original size, err:%s", path, err)
} else {
size = decryptedSize
}
decryptedName, err := d.cipher.DecryptFileName(model.UnwrapObjName(remoteObj).GetName())
if err != nil {
log.Warnf("DecryptFileName failed for %s ,will use original name, err:%s", path, err)
} else {
name = decryptedName
}
} else {
decryptedName, err := d.cipher.DecryptDirName(model.UnwrapObjName(remoteObj).GetName())
if err != nil {
log.Warnf("DecryptDirName failed for %s ,will use original name, err:%s", path, err)
} else {
name = decryptedName
}
}
}
obj := &model.Object{
Path: path,
return &model.Object{
Path: remoteFullPath,
Name: name,
Size: size,
Modified: remoteObj.ModTime(),
IsFolder: remoteObj.IsDir(),
}
return obj, nil
// return nil, errs.ObjectNotFound
Ctime: remoteObj.CreateTime(),
Mask: mask,
}, nil
}
// https://github.com/rclone/rclone/blob/v1.67.0/backend/crypt/cipher.go#L37
const fileHeaderSize = 32
func (d *Crypt) Link(ctx context.Context, file model.Obj, args model.LinkArgs) (*model.Link, error) {
dstDirActualPath, err := d.getActualPathForRemote(file.GetPath(), false)
func (d *Crypt) Link(ctx context.Context, file model.Obj, _ model.LinkArgs) (*model.Link, error) {
remoteStorage, remoteActualPath, err := op.GetStorageAndActualPath(file.GetPath())
if err != nil {
return nil, fmt.Errorf("failed to convert path to remote path: %w", err)
return nil, err
}
remoteLink, remoteFile, err := op.Link(ctx, d.remoteStorage, dstDirActualPath, args)
remoteLink, remoteFile, err := op.Link(ctx, remoteStorage, remoteActualPath, model.LinkArgs{})
if err != nil {
return nil, err
}
@@ -323,30 +311,23 @@ func (d *Crypt) Link(ctx context.Context, file model.Obj, args model.LinkArgs) (
}
func (d *Crypt) MakeDir(ctx context.Context, parentDir model.Obj, dirName string) error {
dstDirActualPath, err := d.getActualPathForRemote(parentDir.GetPath(), true)
remoteStorage, remoteActualPath, err := op.GetStorageAndActualPath(parentDir.GetPath())
if err != nil {
return fmt.Errorf("failed to convert path to remote path: %w", err)
return err
}
dir := d.cipher.EncryptDirName(dirName)
return op.MakeDir(ctx, d.remoteStorage, stdpath.Join(dstDirActualPath, dir))
encryptedName := d.cipher.EncryptDirName(dirName)
return op.MakeDir(ctx, remoteStorage, stdpath.Join(remoteActualPath, encryptedName))
}
func (d *Crypt) Move(ctx context.Context, srcObj, dstDir model.Obj) error {
srcRemoteActualPath, err := d.getActualPathForRemote(srcObj.GetPath(), srcObj.IsDir())
if err != nil {
return fmt.Errorf("failed to convert path to remote path: %w", err)
}
dstRemoteActualPath, err := d.getActualPathForRemote(dstDir.GetPath(), dstDir.IsDir())
if err != nil {
return fmt.Errorf("failed to convert path to remote path: %w", err)
}
return op.Move(ctx, d.remoteStorage, srcRemoteActualPath, dstRemoteActualPath)
_, err := fs.Move(ctx, srcObj.GetPath(), dstDir.GetPath())
return err
}
func (d *Crypt) Rename(ctx context.Context, srcObj model.Obj, newName string) error {
remoteActualPath, err := d.getActualPathForRemote(srcObj.GetPath(), srcObj.IsDir())
remoteStorage, remoteActualPath, err := op.GetStorageAndActualPath(srcObj.GetPath())
if err != nil {
return fmt.Errorf("failed to convert path to remote path: %w", err)
return err
}
var newEncryptedName string
if srcObj.IsDir() {
@@ -354,33 +335,26 @@ func (d *Crypt) Rename(ctx context.Context, srcObj model.Obj, newName string) er
} else {
newEncryptedName = d.cipher.EncryptFileName(newName)
}
return op.Rename(ctx, d.remoteStorage, remoteActualPath, newEncryptedName)
return op.Rename(ctx, remoteStorage, remoteActualPath, newEncryptedName)
}
func (d *Crypt) Copy(ctx context.Context, srcObj, dstDir model.Obj) error {
srcRemoteActualPath, err := d.getActualPathForRemote(srcObj.GetPath(), srcObj.IsDir())
if err != nil {
return fmt.Errorf("failed to convert path to remote path: %w", err)
}
dstRemoteActualPath, err := d.getActualPathForRemote(dstDir.GetPath(), dstDir.IsDir())
if err != nil {
return fmt.Errorf("failed to convert path to remote path: %w", err)
}
return op.Copy(ctx, d.remoteStorage, srcRemoteActualPath, dstRemoteActualPath)
_, err := fs.Copy(ctx, srcObj.GetPath(), dstDir.GetPath())
return err
}
func (d *Crypt) Remove(ctx context.Context, obj model.Obj) error {
remoteActualPath, err := d.getActualPathForRemote(obj.GetPath(), obj.IsDir())
remoteStorage, remoteActualPath, err := op.GetStorageAndActualPath(obj.GetPath())
if err != nil {
return fmt.Errorf("failed to convert path to remote path: %w", err)
return err
}
return op.Remove(ctx, d.remoteStorage, remoteActualPath)
return op.Remove(ctx, remoteStorage, remoteActualPath)
}
func (d *Crypt) Put(ctx context.Context, dstDir model.Obj, streamer model.FileStreamer, up driver.UpdateProgress) error {
dstDirActualPath, err := d.getActualPathForRemote(dstDir.GetPath(), true)
remoteStorage, remoteActualPath, err := op.GetStorageAndActualPath(dstDir.GetPath())
if err != nil {
return fmt.Errorf("failed to convert path to remote path: %w", err)
return err
}
// Encrypt the data into wrappedIn
@@ -404,15 +378,15 @@ func (d *Crypt) Put(ctx context.Context, dstDir model.Obj, streamer model.FileSt
ForceStreamUpload: true,
Exist: streamer.GetExist(),
}
err = op.Put(ctx, d.remoteStorage, dstDirActualPath, streamOut, up, false)
if err != nil {
return err
}
return nil
return op.Put(ctx, remoteStorage, remoteActualPath, streamOut, up)
}
func (d *Crypt) GetDetails(ctx context.Context) (*model.StorageDetails, error) {
remoteDetails, err := op.GetStorageDetails(ctx, d.remoteStorage)
remoteStorage, _, err := op.GetStorageAndActualPath(d.RemotePath)
if err != nil {
return nil, errs.NotImplement
}
remoteDetails, err := op.GetStorageDetails(ctx, remoteStorage)
if err != nil {
return nil, err
}
@@ -421,8 +395,4 @@ func (d *Crypt) GetDetails(ctx context.Context) (*model.StorageDetails, error) {
}, nil
}
//func (d *Safe) Other(ctx context.Context, args model.OtherArgs) (interface{}, error) {
// return nil, errs.NotSupport
//}
var _ driver.Driver = (*Crypt)(nil)
+1 -5
View File
@@ -6,11 +6,6 @@ import (
)
type Addition struct {
// Usually one of two
//driver.RootPath
//driver.RootID
// define other
FileNameEnc string `json:"filename_encryption" type:"select" required:"true" options:"off,standard,obfuscate" default:"off"`
DirNameEnc string `json:"directory_name_encryption" type:"select" required:"true" options:"false,true" default:"false"`
RemotePath string `json:"remote_path" required:"true" help:"This is where the encrypted data stores"`
@@ -32,6 +27,7 @@ var config = driver.Config{
NoCache: true,
DefaultRoot: "/",
NoLinkURL: true,
CheckStatus: true,
}
func init() {
+5 -20
View File
@@ -4,8 +4,6 @@ import (
stdpath "path"
"path/filepath"
"strings"
"github.com/OpenListTeam/OpenList/v4/internal/op"
)
// will give the best guessing based on the path
@@ -15,30 +13,17 @@ func guessPath(path string) (isFolder, secondTry bool) {
return true, false
}
lastSlash := strings.LastIndex(path, "/")
if strings.Index(path[lastSlash:], ".") < 0 {
if !strings.Contains(path[lastSlash:], ".") {
//no dot, try folder then try file
return true, true
}
return false, true
}
func (d *Crypt) getPathForRemote(path string, isFolder bool) (remoteFullPath string) {
if isFolder && !strings.HasSuffix(path, "/") {
path = path + "/"
func (d *Crypt) encryptPath(path string, isFolder bool) string {
if isFolder {
return d.cipher.EncryptDirName(path)
}
dir, fileName := filepath.Split(path)
remoteDir := d.cipher.EncryptDirName(dir)
remoteFileName := ""
if len(strings.TrimSpace(fileName)) > 0 {
remoteFileName = d.cipher.EncryptFileName(fileName)
}
return stdpath.Join(d.RemotePath, remoteDir, remoteFileName)
}
// actual path is used for internal only. any link for user should come from remoteFullPath
func (d *Crypt) getActualPathForRemote(path string, isFolder bool) (string, error) {
_, remoteActualPath, err := op.GetStorageAndActualPath(d.getPathForRemote(path, isFolder))
return remoteActualPath, err
return stdpath.Join(d.cipher.EncryptDirName(dir), d.cipher.EncryptFileName(fileName))
}
+41
View File
@@ -15,6 +15,7 @@ import (
"github.com/OpenListTeam/OpenList/v4/pkg/utils"
"github.com/go-resty/resty/v2"
"github.com/google/uuid"
"golang.org/x/time/rate"
)
type Doubao struct {
@@ -23,6 +24,7 @@ type Doubao struct {
*UploadToken
UserId string
uploadThread int
limiter *rate.Limiter
}
func (d *Doubao) Config() driver.Config {
@@ -61,6 +63,17 @@ func (d *Doubao) Init(ctx context.Context) error {
d.UploadToken = uploadToken
}
if d.LimitRate > 0 {
d.limiter = rate.NewLimiter(rate.Limit(d.LimitRate), 1)
}
return nil
}
func (d *Doubao) WaitLimit(ctx context.Context) error {
if d.limiter != nil {
return d.limiter.Wait(ctx)
}
return nil
}
@@ -69,6 +82,10 @@ func (d *Doubao) Drop(ctx context.Context) error {
}
func (d *Doubao) List(ctx context.Context, dir model.Obj, args model.ListArgs) ([]model.Obj, error) {
if err := d.WaitLimit(ctx); err != nil {
return nil, err
}
var files []model.Obj
fileList, err := d.getFiles(dir.GetID(), "")
if err != nil {
@@ -95,6 +112,10 @@ func (d *Doubao) List(ctx context.Context, dir model.Obj, args model.ListArgs) (
}
func (d *Doubao) Link(ctx context.Context, file model.Obj, args model.LinkArgs) (*model.Link, error) {
if err := d.WaitLimit(ctx); err != nil {
return nil, err
}
var downloadUrl string
if u, ok := file.(*Object); ok {
@@ -160,6 +181,10 @@ func (d *Doubao) Link(ctx context.Context, file model.Obj, args model.LinkArgs)
}
func (d *Doubao) MakeDir(ctx context.Context, parentDir model.Obj, dirName string) error {
if err := d.WaitLimit(ctx); err != nil {
return err
}
var r UploadNodeResp
_, err := d.request("/samantha/aispace/upload_node", http.MethodPost, func(req *resty.Request) {
req.SetBody(base.Json{
@@ -177,6 +202,10 @@ func (d *Doubao) MakeDir(ctx context.Context, parentDir model.Obj, dirName strin
}
func (d *Doubao) Move(ctx context.Context, srcObj, dstDir model.Obj) error {
if err := d.WaitLimit(ctx); err != nil {
return err
}
var r UploadNodeResp
_, err := d.request("/samantha/aispace/move_node", http.MethodPost, func(req *resty.Request) {
req.SetBody(base.Json{
@@ -191,6 +220,10 @@ func (d *Doubao) Move(ctx context.Context, srcObj, dstDir model.Obj) error {
}
func (d *Doubao) Rename(ctx context.Context, srcObj model.Obj, newName string) error {
if err := d.WaitLimit(ctx); err != nil {
return err
}
var r BaseResp
_, err := d.request("/samantha/aispace/rename_node", http.MethodPost, func(req *resty.Request) {
req.SetBody(base.Json{
@@ -207,6 +240,10 @@ func (d *Doubao) Copy(ctx context.Context, srcObj, dstDir model.Obj) (model.Obj,
}
func (d *Doubao) Remove(ctx context.Context, obj model.Obj) error {
if err := d.WaitLimit(ctx); err != nil {
return err
}
var r BaseResp
_, err := d.request("/samantha/aispace/delete_node", http.MethodPost, func(req *resty.Request) {
req.SetBody(base.Json{"node_list": []base.Json{{"id": obj.GetID()}}})
@@ -215,6 +252,10 @@ func (d *Doubao) Remove(ctx context.Context, obj model.Obj) error {
}
func (d *Doubao) Put(ctx context.Context, dstDir model.Obj, file model.FileStreamer, up driver.UpdateProgress) (model.Obj, error) {
if err := d.WaitLimit(ctx); err != nil {
return nil, err
}
// 根据MIME类型确定数据类型
mimetype := file.GetMimetype()
dataType := FileDataType
+9 -4
View File
@@ -10,9 +10,10 @@ type Addition struct {
// driver.RootPath
driver.RootID
// define other
Cookie string `json:"cookie" type:"text"`
UploadThread string `json:"upload_thread" default:"3"`
DownloadApi string `json:"download_api" type:"select" options:"get_file_url,get_download_info" default:"get_file_url"`
Cookie string `json:"cookie" type:"text"`
UploadThread string `json:"upload_thread" default:"3"`
DownloadApi string `json:"download_api" type:"select" options:"get_file_url,get_download_info" default:"get_file_url"`
LimitRate float64 `json:"limit_rate" type:"float" default:"2" help:"limit all api request rate ([limit]r/1s)"`
}
var config = driver.Config{
@@ -23,6 +24,10 @@ var config = driver.Config{
func init() {
op.RegisterDriver(func() driver.Driver {
return &Doubao{}
return &Doubao{
Addition: Addition{
LimitRate: 2,
},
}
})
}
+26 -26
View File
@@ -10,7 +10,6 @@ import (
"fmt"
"hash/crc32"
"io"
"math"
"math/rand"
"net/http"
"net/url"
@@ -18,7 +17,6 @@ import (
"sort"
"strconv"
"strings"
"sync"
"time"
"github.com/OpenListTeam/OpenList/v4/drivers/base"
@@ -62,7 +60,7 @@ const (
VideoDataType = "video"
DefaultChunkSize = int64(5 * 1024 * 1024) // 5MB
MaxRetryAttempts = 3 // 最大重试次数
UserAgent = "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/129.0.0.0 Safari/537.36"
UserAgent = base.UserAgentNT
Region = "cn-north-1"
UploadTimeout = 3 * time.Minute
)
@@ -562,9 +560,7 @@ func (d *Doubao) UploadByMultipart(ctx context.Context, config *UploadConfig, fi
retry.MaxJitter(200*time.Millisecond),
)
var partsMutex sync.Mutex
// 并行上传所有分片
hash := crc32.NewIEEE()
for partIndex := range totalParts {
if utils.IsCanceled(uploadCtx) {
break
@@ -578,32 +574,38 @@ func (d *Doubao) UploadByMultipart(ctx context.Context, config *UploadConfig, fi
size = fileSize - offset
}
var reader io.ReadSeeker
var rateLimitedRd io.Reader
crc32Value := ""
threadG.GoWithLifecycle(errgroup.Lifecycle{
Before: func(ctx context.Context) error {
if reader == nil {
var err error
reader, err = ss.GetSectionReader(offset, size)
if err != nil {
return err
}
hash.Reset()
w, err := utils.CopyWithBuffer(hash, reader)
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 crc32Value == "" {
// 把耗时的计算放在这里,避免阻塞其他协程
crc32Hash := crc32.NewIEEE()
w, err := utils.CopyWithBuffer(crc32Hash, reader)
if w != size {
return fmt.Errorf("failed to read all data: (expect =%d, actual =%d) %w", size, w, err)
}
crc32Value = hex.EncodeToString(hash.Sum(nil))
rateLimitedRd = driver.NewLimitedUploadStream(ctx, reader)
crc32Value = hex.EncodeToString(crc32Hash.Sum(nil))
reader.Seek(0, io.SeekStart)
}
return nil
},
Do: func(ctx context.Context) error {
reader.Seek(0, io.SeekStart)
req, err := http.NewRequestWithContext(ctx, http.MethodPost, fmt.Sprintf("%s?uploadid=%s&part_number=%d&phase=transfer", uploadUrl, uploadID, partNumber), rateLimitedRd)
req, err := http.NewRequestWithContext(
ctx,
http.MethodPost,
uploadUrl,
driver.NewLimitedUploadStream(ctx, reader),
)
if err != nil {
return err
}
query := req.URL.Query()
query.Add("uploadid", uploadID)
query.Add("part_number", strconv.FormatInt(partNumber, 10))
query.Add("phase", "transfer")
req.URL.RawQuery = query.Encode()
req.Header = map[string][]string{
"Referer": {BaseURL + "/"},
"Origin": {BaseURL},
@@ -629,16 +631,14 @@ func (d *Doubao) UploadByMultipart(ctx context.Context, config *UploadConfig, fi
return fmt.Errorf("upload part failed: crc32 mismatch, expected %s, got %s", crc32Value, uploadResp.Data.Crc32)
}
// 记录成功上传的分片
partsMutex.Lock()
parts[partIndex] = UploadPart{
PartNumber: strconv.FormatInt(partNumber, 10),
Etag: uploadResp.Data.Etag,
Crc32: crc32Value,
}
partsMutex.Unlock()
// 更新进度
progress := 10.0 + 90.0*float64(threadG.Success()+1)/float64(totalParts)
up(math.Min(progress, 95.0))
progress := 95 * float64(threadG.Success()+1) / float64(totalParts)
up(progress)
return nil
},
After: func(err error) {
+5 -5
View File
@@ -40,6 +40,7 @@ func (d *DoubaoShare) Drop(ctx context.Context) error {
return nil
}
// 潜在bug:配置二级目录时,可能会出问题
func (d *DoubaoShare) List(ctx context.Context, dir model.Obj, args model.ListArgs) ([]model.Obj, error) {
// 检查是否为根目录
if dir.GetID() == "" && dir.GetPath() == "/" {
@@ -91,18 +92,17 @@ func (d *DoubaoShare) Link(ctx context.Context, file model.Obj, args model.LinkA
downloadUrl = r.Data.OriginalMediaInfo.MainURL
default:
var r GetFileUrlResp
_, err := d.request("/alice/message/get_file_url", http.MethodPost, func(req *resty.Request) {
var r GetDownloadInfoResp
_, err := d.request("/samantha/aispace/get_download_info", http.MethodPost, func(req *resty.Request) {
req.SetBody(base.Json{
"uris": []string{u.Key},
"type": FileNodeType[u.NodeType],
"requests": []base.Json{{"node_id": file.GetID()}},
})
}, &r)
if err != nil {
return nil, err
}
downloadUrl = r.Data.FileUrls[0].MainURL
downloadUrl = r.Data.DownloadInfos[0].MainURL
}
// 生成标准的Content-Disposition
+6 -6
View File
@@ -115,14 +115,14 @@ type FilePath []struct {
UpdateTime int64 `json:"update_time"`
}
type GetFileUrlResp struct {
type GetDownloadInfoResp struct {
BaseResp
Data struct {
FileUrls []struct {
URI string `json:"uri"`
MainURL string `json:"main_url"`
BackURL string `json:"back_url"`
} `json:"file_urls"`
DownloadInfos []struct {
NodeID string `json:"node_id"`
MainURL string `json:"main_url"`
BackupURL string `json:"backup_url"`
} `json:"download_infos"`
} `json:"data"`
}
+1 -1
View File
@@ -43,7 +43,7 @@ const (
FileDataType = "file"
ImgDataType = "image"
VideoDataType = "video"
UserAgent = "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/129.0.0.0 Safari/537.36"
UserAgent = base.UserAgentNT
)
func (d *DoubaoShare) request(path string, method string, callback base.ReqCallback, resp interface{}) ([]byte, error) {
+14 -16
View File
@@ -24,7 +24,6 @@ func (d *Dropbox) refreshToken() error {
ErrorMessage string `json:"text"`
}
_, err := base.RestyClient.R().
SetHeader("User-Agent", "Mozilla/5.0 (Macintosh; Apple macOS 15_5) AppleWebKit/537.36 (KHTML, like Gecko) Safari/537.36 Chrome/138.0.0.0 Openlist/425.6.30").
SetResult(&resp).
SetQueryParams(map[string]string{
"refresh_ui": d.RefreshToken,
@@ -176,12 +175,12 @@ func (d *Dropbox) finishUploadSession(ctx context.Context, toPath string, offset
req.Header.Set("Content-Type", "application/octet-stream")
req.Header.Set("Authorization", "Bearer "+d.AccessToken)
if d.RootNamespaceId != "" {
apiPathRootJson, err := d.buildPathRootHeader()
if err != nil {
return err
apiPathRootJson, err := d.buildPathRootHeader()
if err != nil {
return err
}
req.Header.Set("Dropbox-API-Path-Root", apiPathRootJson)
}
req.Header.Set("Dropbox-API-Path-Root", apiPathRootJson)
}
uploadFinishArgs := UploadFinishArgs{
Commit: struct {
@@ -227,12 +226,12 @@ func (d *Dropbox) startUploadSession(ctx context.Context) (string, error) {
req.Header.Set("Content-Type", "application/octet-stream")
req.Header.Set("Authorization", "Bearer "+d.AccessToken)
if d.RootNamespaceId != "" {
apiPathRootJson, err := d.buildPathRootHeader()
if err != nil {
return "", err
apiPathRootJson, err := d.buildPathRootHeader()
if err != nil {
return "", err
}
req.Header.Set("Dropbox-API-Path-Root", apiPathRootJson)
}
req.Header.Set("Dropbox-API-Path-Root", apiPathRootJson)
}
req.Header.Set("Dropbox-API-Arg", "{\"close\":false}")
res, err := base.HttpClient.Do(req)
@@ -249,9 +248,8 @@ func (d *Dropbox) startUploadSession(ctx context.Context) (string, error) {
}
func (d *Dropbox) buildPathRootHeader() (string, error) {
return utils.Json.MarshalToString(map[string]interface{}{
".tag": "root",
"root": d.RootNamespaceId,
})
return utils.Json.MarshalToString(map[string]interface{}{
".tag": "root",
"root": d.RootNamespaceId,
})
}
+1 -1
View File
@@ -19,7 +19,7 @@ var config = driver.Config{
Name: "FebBox",
NoUpload: true,
DefaultRoot: "0",
LinkCacheType: 1,
LinkCacheMode: driver.LinkCacheIP,
}
func init() {
+1 -3
View File
@@ -113,9 +113,7 @@ func (d *FTP) Link(ctx context.Context, file model.Obj, args model.LinkArgs) (*m
}
return &model.Link{
RangeReader: &model.FileRangeReader{
RangeReaderIF: stream.RateLimitRangeReaderFunc(resultRangeReader),
},
RangeReader: stream.RateLimitRangeReaderFunc(resultRangeReader),
SyncClosers: utils.NewSyncClosers(utils.CloseFunc(conn.Quit)),
}, nil
}
+10 -2
View File
@@ -4,6 +4,7 @@ import (
"context"
"fmt"
"net/http"
stdpath "path"
"strings"
"github.com/OpenListTeam/OpenList/v4/internal/driver"
@@ -51,6 +52,9 @@ func (d *GithubReleases) List(ctx context.Context, dir model.Obj, args model.Lis
if d.Addition.ShowReadme {
files = append(files, point.GetOtherFile(d.GetRequest, args.Refresh)...)
}
if d.Addition.ShowSourceCode {
files = append(files, point.GetSourceCode()...)
}
} else if strings.HasPrefix(point.Point, path) { // 仓库目录的父目录
nextDir := GetNextDir(point.Point, path)
if nextDir == "" {
@@ -67,7 +71,7 @@ func (d *GithubReleases) List(ctx context.Context, dir model.Obj, args model.Lis
}
if !hasSameDir {
files = append(files, File{
Path: path + "/" + nextDir,
Path: stdpath.Join(path, nextDir),
FileName: nextDir,
Size: point.GetLatestSize(),
UpdateAt: point.Release.PublishedAt,
@@ -102,7 +106,7 @@ func (d *GithubReleases) List(ctx context.Context, dir model.Obj, args model.Lis
if !hasSameDir {
files = append(files, File{
FileName: nextDir,
Path: path + "/" + nextDir,
Path: stdpath.Join(path, nextDir),
Size: point.GetAllVersionSize(),
UpdateAt: (*point.Releases)[0].PublishedAt,
CreateAt: (*point.Releases)[0].CreatedAt,
@@ -117,6 +121,10 @@ func (d *GithubReleases) List(ctx context.Context, dir model.Obj, args model.Lis
}
files = append(files, point.GetReleaseByTagName(tagName)...)
if d.Addition.ShowSourceCode {
files = append(files, point.GetSourceCodeByTagName(tagName)...)
}
}
}
}
+2 -1
View File
@@ -6,10 +6,11 @@ import (
)
type Addition struct {
driver.RootID
driver.RootPath
RepoStructure string `json:"repo_structure" type:"text" required:"true" default:"OpenListTeam/OpenList" help:"structure:[path:]org/repo"`
ShowReadme bool `json:"show_readme" type:"bool" default:"true" help:"show README、LICENSE file"`
Token string `json:"token" type:"string" required:"false" help:"GitHub token, if you want to access private repositories or increase the rate limit"`
ShowSourceCode bool `json:"show_source_code" type:"bool" default:"false" help:"show Source code (zip/tar.gz)"`
ShowAllVersion bool `json:"show_all_version" type:"bool" default:"false" help:"show all versions"`
GitHubProxy string `json:"gh_proxy" type:"string" default:"" help:"GitHub proxy, e.g. https://ghproxy.net/github.com or https://gh-proxy.com/github.com "`
}
+60 -5
View File
@@ -2,6 +2,7 @@ package github_releases
import (
"encoding/json"
"path"
"strings"
"time"
@@ -45,10 +46,10 @@ func (m *MountPoint) RequestReleases(get func(url string) (*resty.Response, erro
// 获取最新版本
func (m *MountPoint) GetLatestRelease() []File {
files := make([]File, 0)
files := make([]File, 0, len(m.Release.Assets))
for _, asset := range m.Release.Assets {
files = append(files, File{
Path: m.Point + "/" + asset.Name,
Path: path.Join(m.Point, asset.Name),
FileName: asset.Name,
Size: asset.Size,
Type: "file",
@@ -74,7 +75,7 @@ func (m *MountPoint) GetAllVersion() []File {
files := make([]File, 0)
for _, release := range *m.Releases {
file := File{
Path: m.Point + "/" + release.TagName,
Path: path.Join(m.Point, release.TagName),
FileName: release.TagName,
Size: m.GetSizeByTagName(release.TagName),
Type: "dir",
@@ -97,7 +98,7 @@ func (m *MountPoint) GetReleaseByTagName(tagName string) []File {
files := make([]File, 0)
for _, asset := range item.Assets {
files = append(files, File{
Path: m.Point + "/" + tagName + "/" + asset.Name,
Path: path.Join(m.Point, tagName, asset.Name),
FileName: asset.Name,
Size: asset.Size,
Type: "file",
@@ -143,6 +144,60 @@ func (m *MountPoint) GetAllVersionSize() int64 {
return size
}
func (m *MountPoint) GetSourceCode() []File {
files := make([]File, 0)
// 无法获取文件大小,此处设为 1
files = append(files, File{
Path: path.Join(m.Point, "Source code (zip)"),
FileName: "Source code (zip)",
Size: 1,
Type: "file",
UpdateAt: m.Release.CreatedAt,
CreateAt: m.Release.CreatedAt,
Url: m.Release.ZipballUrl,
})
files = append(files, File{
Path: path.Join(m.Point, "Source code (tar.gz)"),
FileName: "Source code (tar.gz)",
Size: 1,
Type: "file",
UpdateAt: m.Release.CreatedAt,
CreateAt: m.Release.CreatedAt,
Url: m.Release.TarballUrl,
})
return files
}
func (m *MountPoint) GetSourceCodeByTagName(tagName string) []File {
for _, item := range *m.Releases {
if item.TagName == tagName {
files := make([]File, 0)
files = append(files, File{
Path: path.Join(m.Point, "Source code (zip)"),
FileName: "Source code (zip)",
Size: 1,
Type: "file",
UpdateAt: item.CreatedAt,
CreateAt: item.CreatedAt,
Url: item.ZipballUrl,
})
files = append(files, File{
Path: path.Join(m.Point, "Source code (tar.gz)"),
FileName: "Source code (tar.gz)",
Size: 1,
Type: "file",
UpdateAt: item.CreatedAt,
CreateAt: item.CreatedAt,
Url: item.TarballUrl,
})
return files
}
}
return nil
}
func (m *MountPoint) GetOtherFile(get func(url string) (*resty.Response, error), refresh bool) []File {
if m.OtherFile == nil || refresh {
resp, _ := get("https://api.github.com/repos/" + m.Repo + "/contents")
@@ -155,7 +210,7 @@ func (m *MountPoint) GetOtherFile(get func(url string) (*resty.Response, error),
for _, file := range *m.OtherFile {
if strings.HasSuffix(file.Name, ".md") || strings.HasPrefix(file.Name, "LICENSE") {
files = append(files, File{
Path: m.Point + "/" + file.Name,
Path: path.Join(m.Point, file.Name),
FileName: file.Name,
Size: file.Size,
Type: "file",
+81 -2
View File
@@ -27,6 +27,14 @@ import (
// do others that not defined in Driver interface
// Google Drive API field constants
const (
// File list query fields
FilesListFields = "files(id,name,mimeType,size,modifiedTime,createdTime,thumbnailLink,shortcutDetails,md5Checksum,sha1Checksum,sha256Checksum),nextPageToken"
// Single file query fields
FileInfoFields = "id,name,mimeType,size,md5Checksum,sha1Checksum,sha256Checksum"
)
type googleDriveServiceAccount struct {
// Type string `json:"type"`
// ProjectID string `json:"project_id"`
@@ -50,7 +58,6 @@ func (d *GoogleDrive) refreshToken() error {
ErrorMessage string `json:"text"`
}
_, err := base.RestyClient.R().
SetHeader("User-Agent", "Mozilla/5.0 (Macintosh; Apple macOS 15_5) AppleWebKit/537.36 (KHTML, like Gecko) Safari/537.36 Chrome/138.0.0.0 Openlist/425.6.30").
SetResult(&resp).
SetQueryParams(map[string]string{
"refresh_ui": d.RefreshToken,
@@ -235,7 +242,7 @@ func (d *GoogleDrive) getFiles(id string) ([]File, error) {
}
query := map[string]string{
"orderBy": orderBy,
"fields": "files(id,name,mimeType,size,modifiedTime,createdTime,thumbnailLink,shortcutDetails,md5Checksum,sha1Checksum,sha256Checksum),nextPageToken",
"fields": FilesListFields,
"pageSize": "1000",
"q": fmt.Sprintf("'%s' in parents and trashed = false", id),
//"includeItemsFromAllDrives": "true",
@@ -249,11 +256,82 @@ func (d *GoogleDrive) getFiles(id string) ([]File, error) {
return nil, err
}
pageToken = resp.NextPageToken
// Batch process shortcuts, API calls only for file shortcuts
shortcutTargetIds := make([]string, 0)
shortcutIndices := make([]int, 0)
// Collect target IDs of all file shortcuts (skip folder shortcuts)
for i := range resp.Files {
if resp.Files[i].MimeType == "application/vnd.google-apps.shortcut" &&
resp.Files[i].ShortcutDetails.TargetId != "" &&
resp.Files[i].ShortcutDetails.TargetMimeType != "application/vnd.google-apps.folder" {
shortcutTargetIds = append(shortcutTargetIds, resp.Files[i].ShortcutDetails.TargetId)
shortcutIndices = append(shortcutIndices, i)
}
}
// Batch get target file info (only for file shortcuts)
if len(shortcutTargetIds) > 0 {
targetFiles := d.batchGetTargetFilesInfo(shortcutTargetIds)
// Update shortcut file info
for j, targetId := range shortcutTargetIds {
if targetFile, exists := targetFiles[targetId]; exists {
fileIndex := shortcutIndices[j]
if targetFile.Size != "" {
resp.Files[fileIndex].Size = targetFile.Size
}
if targetFile.MD5Checksum != "" {
resp.Files[fileIndex].MD5Checksum = targetFile.MD5Checksum
}
if targetFile.SHA1Checksum != "" {
resp.Files[fileIndex].SHA1Checksum = targetFile.SHA1Checksum
}
if targetFile.SHA256Checksum != "" {
resp.Files[fileIndex].SHA256Checksum = targetFile.SHA256Checksum
}
}
}
}
res = append(res, resp.Files...)
}
return res, nil
}
// getTargetFileInfo gets target file details for shortcuts
func (d *GoogleDrive) getTargetFileInfo(targetId string) (File, error) {
var targetFile File
url := fmt.Sprintf("https://www.googleapis.com/drive/v3/files/%s", targetId)
query := map[string]string{
"fields": FileInfoFields,
}
_, err := d.request(url, http.MethodGet, func(req *resty.Request) {
req.SetQueryParams(query)
}, &targetFile)
if err != nil {
return File{}, err
}
return targetFile, nil
}
// batchGetTargetFilesInfo batch gets target file info, sequential processing to avoid concurrency complexity
func (d *GoogleDrive) batchGetTargetFilesInfo(targetIds []string) map[string]File {
if len(targetIds) == 0 {
return make(map[string]File)
}
result := make(map[string]File)
// Sequential processing to avoid concurrency complexity
for _, targetId := range targetIds {
file, err := d.getTargetFileInfo(targetId)
if err == nil {
result[targetId] = file
}
}
return result
}
func (d *GoogleDrive) chunkUpload(ctx context.Context, file model.FileStreamer, url string, up driver.UpdateProgress) error {
defaultChunkSize := d.ChunkSize * 1024 * 1024
ss, err := stream.NewStreamSectionReader(file, int(defaultChunkSize), &up)
@@ -304,6 +382,7 @@ func (d *GoogleDrive) chunkUpload(ctx context.Context, file model.FileStreamer,
up(float64(offset+chunkSize) / float64(file.GetSize()) * 100)
return nil
},
retry.Context(ctx),
retry.Attempts(3),
retry.DelayType(retry.BackOffDelay),
retry.Delay(time.Second))
+1 -2
View File
@@ -2,7 +2,6 @@ package halalcloudopen
import (
"context"
"strconv"
"github.com/OpenListTeam/OpenList/v4/internal/model"
sdkModel "github.com/halalcloud/golang-sdk-lite/halalcloud/model"
@@ -19,7 +18,7 @@ func (d *HalalCloudOpen) getFiles(ctx context.Context, dir model.Obj) ([]model.O
result, err := d.sdkUserFileService.List(ctx, &sdkUserFile.FileListRequest{
Parent: &sdkUserFile.File{Path: dir.GetPath()},
ListInfo: &sdkModel.ScanListRequest{
Limit: strconv.FormatInt(limit, 10),
Limit: limit,
Token: token,
},
})
+36 -36
View File
@@ -16,6 +16,7 @@ import (
"github.com/OpenListTeam/OpenList/v4/internal/model"
sdkUserFile "github.com/halalcloud/golang-sdk-lite/halalcloud/services/userfile"
"github.com/ipfs/go-cid"
log "github.com/sirupsen/logrus"
)
func (d *HalalCloudOpen) put(ctx context.Context, dstDir model.Obj, fileStream model.FileStreamer, up driver.UpdateProgress) (model.Obj, error) {
@@ -51,52 +52,39 @@ func (d *HalalCloudOpen) put(ctx context.Context, dstDir model.Obj, fileStream m
Version: 1,
}
blockSize := uploadTask.BlockSize
useSingleUpload := true
//
if fileStream.GetSize() <= int64(blockSize) || d.uploadThread <= 1 {
useSingleUpload = true
}
// Not sure whether FileStream supports concurrent read and write operations, so currently using single-threaded upload to ensure safety.
// read file
if useSingleUpload {
bufferSize := int(blockSize)
buffer := make([]byte, bufferSize)
reader := driver.NewLimitedUploadStream(ctx, fileStream)
teeReader := io.TeeReader(reader, driver.NewProgress(fileStream.GetSize(), up))
// fileStream.Seek(0, os.SEEK_SET)
for {
n, err := teeReader.Read(buffer)
if n > 0 {
data := buffer[:n]
uploadCid, err := postFileSlice(ctx, data, uploadTask.Task, uploadTask.UploadAddress, prefix, retryTimes)
bufferSize := int(blockSize)
buffer := make([]byte, bufferSize)
offset := 0
teeReader := io.TeeReader(fileStream, driver.NewProgress(fileStream.GetSize(), up))
for {
n, err := teeReader.Read(buffer[offset:]) // 这里 len(buf[offset:]) <= 4MB
if n > 0 {
offset += n
if offset == int(blockSize) {
uploadCid, err := postFileSlice(ctx, buffer, uploadTask.Task, uploadTask.UploadAddress, prefix, retryTimes)
if err != nil {
return nil, err
}
slicesList = append(slicesList, uploadCid.String())
}
if err == io.EOF || n == 0 {
break
offset = 0
}
}
} else {
// TODO: implement multipart upload, currently using single-threaded upload to ensure safety.
bufferSize := int(blockSize)
buffer := make([]byte, bufferSize)
reader := driver.NewLimitedUploadStream(ctx, fileStream)
teeReader := io.TeeReader(reader, driver.NewProgress(fileStream.GetSize(), up))
for {
n, err := teeReader.Read(buffer)
if n > 0 {
data := buffer[:n]
uploadCid, err := postFileSlice(ctx, data, uploadTask.Task, uploadTask.UploadAddress, prefix, retryTimes)
if err != nil {
return nil, err
if err != nil {
if err == io.EOF {
if offset > 0 {
uploadCid, err := postFileSlice(ctx, buffer[:offset], uploadTask.Task, uploadTask.UploadAddress, prefix, retryTimes)
if err != nil {
return nil, err
}
slicesList = append(slicesList, uploadCid.String())
}
slicesList = append(slicesList, uploadCid.String())
}
if err == io.EOF || n == 0 {
break
}
return nil, err
}
}
newFile, err := makeFile(ctx, slicesList, uploadTask.Task, uploadTask.UploadAddress, retryTimes)
@@ -118,6 +106,7 @@ func makeFile(ctx context.Context, fileSlice []string, taskID string, uploadAddr
if ctx.Err() != nil {
return nil, err
}
log.Errorf("make file slice failed, retrying... error: %s", err.Error())
if strings.Contains(err.Error(), "not found") {
return nil, err
}
@@ -156,15 +145,23 @@ func doMakeFile(fileSlice []string, taskID string, uploadAddress string) (*sdkUs
if httpResponse.StatusCode != http.StatusOK && httpResponse.StatusCode != http.StatusCreated {
b, _ := io.ReadAll(httpResponse.Body)
message := string(b)
log.Errorf("make file failed, status code: %d, message: %s", httpResponse.StatusCode, message)
return nil, fmt.Errorf("mk file slice failed, status code: %d, message: %s", httpResponse.StatusCode, message)
}
b, _ := io.ReadAll(httpResponse.Body)
var result *sdkUserFile.File
var result *UploadedFile
err = json.Unmarshal(b, &result)
if err != nil {
log.Errorf("make file failed from response, status code: %d, message: %s", httpResponse.StatusCode, string(b))
return nil, err
}
return result, nil
return &sdkUserFile.File{
Identity: result.Identity,
Path: result.Path,
Size: result.Size,
ContentIdentity: result.ContentIdentity,
}, nil
}
func postFileSlice(ctx context.Context, fileSlice []byte, taskID string, uploadAddress string, preix cid.Prefix, retry int) (cid.Cid, error) {
var lastError error = nil
@@ -214,9 +211,11 @@ func doPostFileSlice(fileSlice []byte, taskID string, uploadAddress string, prei
}
httpResponse, err := httpClient.Do(&httpRequest)
if err != nil {
log.Errorf("access %s failed, method: %s", accessUrl, http.MethodGet)
return cid.Undef, err
}
if httpResponse.StatusCode != http.StatusOK {
log.Errorf("access %s failed, method: %s, status code: %d", accessUrl, http.MethodGet, httpResponse.StatusCode)
return cid.Undef, fmt.Errorf("upload file slice failed, status code: %d", httpResponse.StatusCode)
}
var result bool
@@ -250,6 +249,7 @@ func doPostFileSlice(fileSlice []byte, taskID string, uploadAddress string, prei
if httpResponse.StatusCode != http.StatusOK && httpResponse.StatusCode != http.StatusCreated {
b, _ := io.ReadAll(httpResponse.Body)
message := string(b)
log.Errorf("upload file slice failed, status code: %d, message: %s", httpResponse.StatusCode, message)
return cid.Undef, fmt.Errorf("upload file slice failed, status code: %d, message: %s", httpResponse.StatusCode, message)
}
//
+8
View File
@@ -30,3 +30,11 @@ func init() {
return &HalalCloudOpen{}
})
}
type UploadedFile struct {
Identity string `json:"identity"`
UserIdentity string `json:"user_identity"`
Path string `json:"path"`
Size int64 `json:"size"`
ContentIdentity string `json:"content_identity"`
}
+3 -3
View File
@@ -152,8 +152,7 @@ func (d *ILanZou) Link(ctx context.Context, file model.Obj, args model.LinkArgs)
req := base.NoRedirectClient.R()
req.SetHeaders(map[string]string{
"Referer": d.conf.site + "/",
"User-Agent": "Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/125.0.0.0 Safari/537.36 Edg/125.0.0.0",
"Referer": d.conf.site + "/",
})
if d.Addition.Ip != "" {
req.SetHeader("X-Forwarded-For", d.Addition.Ip)
@@ -409,9 +408,10 @@ func (d *ILanZou) GetDetails(ctx context.Context) (*model.StorageDetails, error)
if err != nil {
return nil, err
}
vipSize := utils.Json.Get(res, "map", "vipSize").ToUint64() * 1024
totalSize := utils.Json.Get(res, "map", "totalSize").ToUint64() * 1024
rewardSize := utils.Json.Get(res, "map", "rewardSize").ToUint64() * 1024
total := totalSize + rewardSize
total := totalSize + rewardSize + vipSize
used := utils.Json.Get(res, "map", "usedSize").ToUint64() * 1024
return &model.StorageDetails{
DiskUsage: driver.DiskUsageFromUsedAndTotal(used, total),
+2 -3
View File
@@ -71,9 +71,8 @@ func (d *ILanZou) request(pathname, method string, callback base.ReqCallback, pr
req.SetHeaders(map[string]string{
"Origin": d.conf.site,
"Referer": d.conf.site + "/",
"User-Agent": "Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/125.0.0.0 Safari/537.36 Edg/125.0.0.0",
"Accept-Encoding": "gzip, deflate, br, zstd",
"Accept-Language": "zh-CN,zh;q=0.9,en;q=0.8,en-GB;q=0.7,en-US;q=0.6,mt;q=0.5",
"Accept-Encoding": "gzip",
"Accept-Language": "zh-CN,zh;q=0.9,en-US,en;q=0.8",
})
if d.Addition.Ip != "" {
+1 -1
View File
@@ -31,7 +31,7 @@ func (d *LanZou) GetAddition() driver.Additional {
func (d *LanZou) Init(ctx context.Context) (err error) {
if d.UserAgent == "" {
d.UserAgent = "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.39 (KHTML, like Gecko) Chrome/89.0.4389.111 Safari/537.39"
d.UserAgent = base.UserAgentNT
}
switch d.Type {
case "account":
+1 -1
View File
@@ -17,7 +17,7 @@ type Addition struct {
SharePassword string `json:"share_password"`
BaseUrl string `json:"baseUrl" required:"true" default:"https://pc.woozooo.com" help:"basic URL for file operation"`
ShareUrl string `json:"shareUrl" required:"true" default:"https://pan.lanzoui.com" help:"used to get the sharing page"`
UserAgent string `json:"user_agent" required:"true" default:"Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.39 (KHTML, like Gecko) Chrome/89.0.4389.111 Safari/537.39"`
UserAgent string `json:"user_agent" required:"true" default:"Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.39 (KHTML, like Gecko) Chrome/142.0.0.0 Safari/537.39"`
RepairFileInfo bool `json:"repair_file_info" help:"To use webdav, you need to enable it"`
}
+4 -9
View File
@@ -2,6 +2,7 @@ package LenovoNasShare
import (
"context"
"fmt"
"net/http"
"net/url"
"strings"
@@ -47,12 +48,7 @@ func (d *LenovoNasShare) Drop(ctx context.Context) error {
func (d *LenovoNasShare) List(ctx context.Context, dir model.Obj, args model.ListArgs) ([]model.Obj, error) {
d.checkStoken() // 检查stoken是否过期
files := make([]File, 0)
path := dir.GetPath()
if path == "" && !d.ShowRootFolder && d.RootFolderPath != "" {
path = d.RootFolderPath
}
path := fmt.Sprintf("/%s", strings.Trim(dir.GetPath(), "/"))
var resp Files
query := map[string]string{
@@ -69,15 +65,14 @@ func (d *LenovoNasShare) List(ctx context.Context, dir model.Obj, args model.Lis
return nil, err
}
files = append(files, resp.Data.List...)
return utils.SliceConvert(files, func(src File) (model.Obj, error) {
return utils.SliceConvert(resp.Data.List, func(src File) (model.Obj, error) {
if src.IsDir() {
return src, nil
}
return &model.ObjThumb{
Object: model.Object{
Name: src.GetName(),
Path: src.GetPath(),
Size: src.GetSize(),
Modified: src.ModTime(),
IsFolder: src.IsDir(),
+16
View File
@@ -0,0 +1,16 @@
//go:build !windows && !plan9 && !netbsd && !aix && !illumos && !solaris && !js
package local
import (
"os"
"path/filepath"
"syscall"
)
func copyNamedPipe(dstPath string, mode os.FileMode, dirMode os.FileMode) error {
if err := os.MkdirAll(filepath.Dir(dstPath), dirMode); err != nil {
return err
}
return syscall.Mkfifo(dstPath, uint32(mode))
}
+9
View File
@@ -0,0 +1,9 @@
//go:build windows || plan9 || netbsd || aix || illumos || solaris || js
package local
import "os"
func copyNamedPipe(_ string, _, _ os.FileMode) error {
return nil
}
+8 -17
View File
@@ -23,7 +23,6 @@ import (
"github.com/OpenListTeam/OpenList/v4/pkg/utils"
"github.com/OpenListTeam/OpenList/v4/server/common"
"github.com/OpenListTeam/times"
cp "github.com/otiai10/copy"
log "github.com/sirupsen/logrus"
_ "golang.org/x/image/webp"
)
@@ -297,16 +296,9 @@ func (d *Local) Move(ctx context.Context, srcObj, dstDir model.Obj) error {
return fmt.Errorf("the destination folder is a subfolder of the source folder")
}
err := os.Rename(srcPath, dstPath)
if err != nil && strings.Contains(err.Error(), "invalid cross-device link") {
// 跨设备移动,先复制再删除
if err := d.Copy(ctx, srcObj, dstDir); err != nil {
return err
}
// 复制成功后直接删除源文件/文件夹
if srcObj.IsDir() {
return os.RemoveAll(srcObj.GetPath())
}
return os.Remove(srcObj.GetPath())
if isCrossDeviceError(err) {
// 跨设备移动,变更为移动任务
return errs.NotImplement
}
if err == nil {
srcParent := filepath.Dir(srcPath)
@@ -347,15 +339,14 @@ func (d *Local) Copy(_ context.Context, srcObj, dstDir model.Obj) error {
if utils.IsSubPath(srcPath, dstPath) {
return fmt.Errorf("the destination folder is a subfolder of the source folder")
}
// Copy using otiai10/copy to perform more secure & efficient copy
err := cp.Copy(srcPath, dstPath, cp.Options{
Sync: true, // Sync file to disk after copy, may have performance penalty in filesystem such as ZFS
PreserveTimes: true,
PreserveOwner: true,
})
info, err := os.Lstat(srcPath)
if err != nil {
return err
}
// 复制regular文件会返回errs.NotImplement, 转为复制任务
if err = d.tryCopy(srcPath, dstPath, info); err != nil {
return err
}
if d.directoryMap.Has(filepath.Dir(dstPath)) {
d.directoryMap.UpdateDirSize(filepath.Dir(dstPath))
+80 -1
View File
@@ -3,6 +3,7 @@ package local
import (
"bytes"
"encoding/json"
"errors"
"fmt"
"io/fs"
"os"
@@ -14,7 +15,9 @@ import (
"strings"
"sync"
"github.com/KarpelesLab/reflink"
"github.com/OpenListTeam/OpenList/v4/internal/conf"
"github.com/OpenListTeam/OpenList/v4/internal/errs"
"github.com/OpenListTeam/OpenList/v4/internal/model"
"github.com/OpenListTeam/OpenList/v4/pkg/utils"
"github.com/disintegration/imaging"
@@ -148,7 +151,7 @@ func (d *Local) getThumb(file model.Obj) (*bytes.Buffer, *string, error) {
return nil, nil, err
}
if d.ThumbCacheFolder != "" {
err = os.WriteFile(filepath.Join(d.ThumbCacheFolder, thumbName), buf.Bytes(), 0666)
err = os.WriteFile(filepath.Join(d.ThumbCacheFolder, thumbName), buf.Bytes(), 0o666)
if err != nil {
return nil, nil, err
}
@@ -405,3 +408,79 @@ func (m *DirectoryMap) DeleteDirNode(dirname string) error {
return nil
}
func (d *Local) tryCopy(srcPath, dstPath string, info os.FileInfo) error {
if info.Mode()&os.ModeDevice != 0 {
return errors.New("cannot copy a device")
} else if info.Mode()&os.ModeSymlink != 0 {
return d.copySymlink(srcPath, dstPath)
} else if info.Mode()&os.ModeNamedPipe != 0 {
return copyNamedPipe(dstPath, info.Mode(), os.FileMode(d.mkdirPerm))
} else if info.IsDir() {
return d.recurAndTryCopy(srcPath, dstPath)
} else {
return tryReflinkCopy(srcPath, dstPath)
}
}
func (d *Local) copySymlink(srcPath, dstPath string) error {
linkOrig, err := os.Readlink(srcPath)
if err != nil {
return err
}
dstDir := filepath.Dir(dstPath)
if !filepath.IsAbs(linkOrig) {
srcDir := filepath.Dir(srcPath)
rel, err := filepath.Rel(dstDir, srcDir)
if err != nil {
rel, err = filepath.Abs(srcDir)
}
if err != nil {
return err
}
linkOrig = filepath.Clean(filepath.Join(rel, linkOrig))
}
err = os.MkdirAll(dstDir, os.FileMode(d.mkdirPerm))
if err != nil {
return err
}
return os.Symlink(linkOrig, dstPath)
}
func (d *Local) recurAndTryCopy(srcPath, dstPath string) error {
err := os.MkdirAll(dstPath, os.FileMode(d.mkdirPerm))
if err != nil {
return err
}
files, err := readDir(srcPath)
if err != nil {
return err
}
for _, f := range files {
if !f.IsDir() {
sp := filepath.Join(srcPath, f.Name())
dp := filepath.Join(dstPath, f.Name())
if err = d.tryCopy(sp, dp, f); err != nil {
return err
}
}
}
for _, f := range files {
if f.IsDir() {
sp := filepath.Join(srcPath, f.Name())
dp := filepath.Join(dstPath, f.Name())
if err = d.recurAndTryCopy(sp, dp); err != nil {
return err
}
}
}
return nil
}
func tryReflinkCopy(srcPath, dstPath string) error {
err := reflink.Always(srcPath, dstPath)
if errors.Is(err, reflink.ErrReflinkUnsupported) || errors.Is(err, reflink.ErrReflinkFailed) || isCrossDeviceError(err) {
return errs.NotImplement
}
return err
}
+6
View File
@@ -3,11 +3,13 @@
package local
import (
"errors"
"io/fs"
"strings"
"syscall"
"github.com/OpenListTeam/OpenList/v4/internal/model"
"golang.org/x/sys/unix"
)
func isHidden(f fs.FileInfo, _ string) bool {
@@ -27,3 +29,7 @@ func getDiskUsage(path string) (model.DiskUsage, error) {
FreeSpace: free,
}, nil
}
func isCrossDeviceError(err error) bool {
return errors.Is(err, unix.EXDEV)
}
+4
View File
@@ -49,3 +49,7 @@ func getDiskUsage(path string) (model.DiskUsage, error) {
FreeSpace: freeBytes,
}, nil
}
func isCrossDeviceError(err error) bool {
return errors.Is(err, windows.ERROR_NOT_SAME_DEVICE)
}
+11 -6
View File
@@ -60,20 +60,24 @@ func (d *Mediafire) GetAddition() driver.Additional {
// Init initializes the MediaFire driver with session token and cookie validation
func (d *Mediafire) Init(ctx context.Context) error {
if d.SessionToken == "" {
return fmt.Errorf("Init :: [MediaFire] {critical} missing sessionToken")
}
if d.Cookie == "" {
return fmt.Errorf("Init :: [MediaFire] {critical} missing Cookie")
}
// If SessionToken is empty, try to get it from cookie
if d.SessionToken == "" {
if _, err := d.getSessionToken(ctx); err != nil {
return fmt.Errorf("Init :: [MediaFire] {critical} failed to get session token from cookie: %w", err)
}
}
// Setup rate limiter if rate limit is configured
if d.LimitRate > 0 {
d.limiter = rate.NewLimiter(rate.Limit(d.LimitRate), 1)
}
// Validate and refresh session token if needed
if _, err := d.getSessionToken(ctx); err != nil {
d.renewToken(ctx)
// Avoids 10 mins token expiry (6- 9)
@@ -387,8 +391,8 @@ func (d *Mediafire) Put(ctx context.Context, dstDir model.Obj, file model.FileSt
}
} else {
pollKey = checkResp.Response.ResumableUpload.UploadKey
up(100.0)
}
defer up(100.0)
pollResp, err := d.pollUpload(ctx, pollKey)
if err != nil {
@@ -429,3 +433,4 @@ func (d *Mediafire) GetDetails(ctx context.Context) (*model.StorageDetails, erro
}
var _ driver.Driver = (*Mediafire)(nil)

Some files were not shown because too many files have changed in this diff Show More