From 4c161b63c65be763d8a014f9c793647390fe079a Mon Sep 17 00:00:00 2001 From: admin <362324317@qq.com> Date: Fri, 22 May 2026 19:02:46 +0800 Subject: [PATCH] feat: integrate Quark and UC drive APIs Add optional drive API capabilities for Quark and UC adapters, including directory lookup/creation, rename, move, delete, task polling, Quark recycle cleanup, and UC share staging folder support.\n\nAdd unittest coverage for capability declarations, production transfer path save_dir resolution, staging-folder flow, delete_files behavior, task polling params, and HTTP/JSON error handling.\n\nDocument the Hong Kong test server deployment boundary and verification commands. --- cloudsearch_transfer/adapter/base.py | 74 ++++- .../adapter/quark/__init__.py | 197 +++++++++-- cloudsearch_transfer/adapter/uc/__init__.py | 243 +++++++++++--- .../test_quark_uc_capabilities_unittest.py | 308 ++++++++++++++++++ docs/deployment/test-server.md | 19 ++ 5 files changed, 768 insertions(+), 73 deletions(-) create mode 100644 cloudsearch_transfer/tests/test_quark_uc_capabilities_unittest.py create mode 100644 docs/deployment/test-server.md diff --git a/cloudsearch_transfer/adapter/base.py b/cloudsearch_transfer/adapter/base.py index 08bac8a..88c6abb 100644 --- a/cloudsearch_transfer/adapter/base.py +++ b/cloudsearch_transfer/adapter/base.py @@ -76,6 +76,18 @@ class BaseCloudDriveAdapter(ABC): # URL匹配正则(子类覆盖) URL_PATTERNS: List[str] = [] + # 可选 Drive API 能力;子类按需覆盖为 True 并实现对应方法。 + capabilities: Dict[str, bool] = { + "ensure_dir": False, + "save_files": False, + "poll_task": False, + "rename": False, + "move_files": False, + "delete_files": False, + "cleanup_recycle": False, + "share_staging_folder": False, + } + # 默认请求头 DEFAULT_HEADERS: Dict[str, str] = { "User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) " @@ -235,6 +247,40 @@ class BaseCloudDriveAdapter(ABC): """广告过滤(默认不实现,子类可覆盖)""" return file_ids + + # ─── Optional Drive API capability protocol ───────────────────── + + def _unsupported_capability(self, capability: str) -> None: + raise TransferError( + TransferErrorCode.NETWORK_ERROR, + message=f"{self.PLATFORM_KEY or self.PLATFORM_NAME} 不支持 Drive API 能力: {capability}", + platform=self.PLATFORM_KEY, + ) + + def ensure_dir(self, dir_path: str) -> str: + self._unsupported_capability("ensure_dir") + + def get_fids(self, file_paths: List[str]) -> List[Dict[str, Any]]: + self._unsupported_capability("get_fids") + + def mkdir(self, dir_path: str) -> Dict[str, Any]: + self._unsupported_capability("mkdir") + + def rename(self, fid: str, file_name: str) -> Dict[str, Any]: + self._unsupported_capability("rename") + + def move_files(self, fids: List[str], to_pdir_fid: str) -> Dict[str, Any]: + self._unsupported_capability("move_files") + + def delete_files(self, fids: List[str]) -> Dict[str, Any]: + self._unsupported_capability("delete_files") + + def query_task(self, task_id: str) -> Dict[str, Any]: + self._unsupported_capability("poll_task") + + def cleanup_recycle(self, fids: List[str]) -> Dict[str, Any]: + self._unsupported_capability("cleanup_recycle") + # ─── HTTP 工具方法 ───────────────────────────────────── def _get(self, url: str, params: dict = None, headers: dict = None, @@ -271,6 +317,28 @@ class BaseCloudDriveAdapter(ABC): raise TransferError(TransferErrorCode.NETWORK_ERROR, message=str(last_exc), platform=self.PLATFORM_KEY) + + def _drive_api_json(self, resp: requests.Response, context: str = "网盘 API") -> Dict[str, Any]: + """Validate HTTP response and decode JSON for drive helper APIs.""" + try: + resp.raise_for_status() + except requests.HTTPError as exc: + raise TransferError( + TransferErrorCode.NETWORK_ERROR, + message=f"{context} HTTP错误: {exc}", + platform=self.PLATFORM_KEY, + details={"status_code": getattr(resp, "status_code", None)}, + ) from exc + try: + return resp.json() + except ValueError as exc: + text = getattr(resp, "text", "") or "" + raise TransferError( + TransferErrorCode.NETWORK_ERROR, + message=f"{context} JSON解析失败: {text[:200]}", + platform=self.PLATFORM_KEY, + ) from exc + def _poll_task(self, task_url: str, task_id: str, status_field: str = "status", success_value: Any = 2, @@ -292,10 +360,12 @@ class BaseCloudDriveAdapter(ABC): details={"task_id": task_id}) try: - params = query_params or {} + base_params = query_params(attempt) if callable(query_params) else (query_params or {}) + params = dict(base_params) params["task_id"] = task_id resp = self._get(task_url, params=params, retry=1) - data = resp.json().get("data", resp.json()) + payload = resp.json() + data = payload.get("data", payload) current_status = data.get(status_field) if current_status == success_value: diff --git a/cloudsearch_transfer/adapter/quark/__init__.py b/cloudsearch_transfer/adapter/quark/__init__.py index 55d8d93..3f72ae3 100644 --- a/cloudsearch_transfer/adapter/quark/__init__.py +++ b/cloudsearch_transfer/adapter/quark/__init__.py @@ -55,19 +55,21 @@ class QuarkAdapter(BaseCloudDriveAdapter): r"pan\.quark\.cn/s/(\w+)", ] + + capabilities: Dict[str, bool] = { + "ensure_dir": True, + "save_files": True, + "poll_task": True, + "rename": True, + "move_files": True, + "delete_files": True, + "cleanup_recycle": True, + "share_staging_folder": False, + } + def __init__(self, config: PlatformConfig, transfer_config: TransferConfig) -> None: - """初始化夸克适配器。 - - Args: - config: 平台配置(含 Cookie 等)。 - transfer_config: 全局转存配置(超时、重试、轮询参数等)。 - """ - super().__init__(config, transfer_config) - - # 初始化三个子模块 - self._credential: QuarkCredentialManager = QuarkCredentialManager( - cookie=config.cookie - ) + """初始化适配器。""" + self._credential: QuarkCredentialManager = QuarkCredentialManager(cookie=config.cookie) self._transfer_engine: QuarkTransfer = QuarkTransfer( credential=self._credential, timeout=transfer_config.request_timeout, @@ -78,6 +80,7 @@ class QuarkAdapter(BaseCloudDriveAdapter): credential=self._credential, timeout=transfer_config.request_timeout, ) + super().__init__(config, transfer_config) # ═══════════════════════════════════════════════════════════════ # 公开接口实现 @@ -116,8 +119,9 @@ class QuarkAdapter(BaseCloudDriveAdapter): platform=self.PLATFORM_KEY, ) - # 目标目录:默认根目录 "0" - target_dir: str = save_dir or self.config.save_dir or "0" + # 目标目录:支持 fid 或路径;路径会先创建/解析为 fid。 + requested_dir: str = save_dir or self.config.save_dir or "/" + target_dir: str = requested_dir if requested_dir and not str(requested_dir).startswith("/") else self.ensure_dir(requested_dir) # 分享密码 pwd: str = share_password or self.config.share_password or "" @@ -250,18 +254,13 @@ class QuarkAdapter(BaseCloudDriveAdapter): def _save_files(self, pwd_id: str, detail: dict, save_dir: str) -> List[str]: """转存文件到自己的夸克网盘(基类 transfer() 流程中的步骤③④)。 - Args: - pwd_id: 分享 ID。 - detail: 分享详情(来自 _get_share_detail)。 - save_dir: 目标目录 ID。 - - Returns: - 转存后的新文件 ID 列表。 + save_dir 可传 fid 或路径;路径会先通过 ensure_dir 创建/解析为 fid。 """ # 需要 stoken,从 detail 间接获取(重新请求) stoken: str = self._transfer_engine._get_stoken(pwd_id) + target_fid = save_dir if save_dir and not str(save_dir).startswith("/") else self.ensure_dir(save_dir or "/") task_id: str = self._transfer_engine._init_save( - pwd_id, stoken, detail, to_pdir_fid=save_dir + pwd_id, stoken, detail, to_pdir_fid=target_fid ) return self._transfer_engine._poll_save_task(task_id) @@ -351,6 +350,160 @@ class QuarkAdapter(BaseCloudDriveAdapter): logger.warning("[QuarkAdapter] Cannot fetch file list for ad filtering, skipping") return file_ids + + # ─── Drive API capability helpers ───────────────────────────── + + @staticmethod + def _normalize_dir_path(dir_path: str) -> str: + path = "/" + str(dir_path or "").strip().strip("/") + return "/" if path == "/" else path + + @staticmethod + def _api_success(payload: Dict[str, Any]) -> bool: + if not isinstance(payload, dict): + return False + code = payload.get("code") + status = payload.get("status") + return code == 0 or status in (0, 200) + + def ensure_dir(self, dir_path: str) -> str: + normalized = self._normalize_dir_path(dir_path) + if normalized == "/": + return "0" + + parts = [part for part in normalized.strip("/").split("/") if part] + prefixes = ["/" + "/".join(parts[:idx]) for idx in range(1, len(parts) + 1)] + existing = { + item.get("file_path"): str(item.get("fid")) + for item in self.get_fids(prefixes) + if item.get("file_path") and item.get("fid") + } + + leaf_fid = existing.get(normalized) + if leaf_fid: + return leaf_fid + + for prefix in prefixes: + if prefix in existing: + continue + result = self.mkdir(prefix) + if self._api_success(result) and result.get("data", {}).get("fid"): + existing[prefix] = str(result["data"]["fid"]) + continue + raise TransferError( + TransferErrorCode.NETWORK_ERROR, + message=f"创建目录失败: {result.get('message', result)}", + platform=self.PLATFORM_KEY, + ) + return existing[normalized] + + def mkdir(self, dir_path: str) -> Dict[str, Any]: + url = "https://drive-pc.quark.cn/1/clouddrive/file" + params = {"pr": "ucpro", "fr": "pc", "uc_param_str": ""} + payload = { + "pdir_fid": "0", + "file_name": "", + "dir_path": self._normalize_dir_path(dir_path), + "dir_init_lock": False, + } + return self._drive_api_json(self._post(url, json_data=payload, params=params, headers=self._credential.get_headers()), context="写入网盘目录") + + def rename(self, fid: str, file_name: str) -> Dict[str, Any]: + url = "https://drive-pc.quark.cn/1/clouddrive/file/rename" + params = {"pr": "ucpro", "fr": "pc", "uc_param_str": ""} + payload = {"fid": fid, "file_name": file_name} + return self._drive_api_json(self._post(url, json_data=payload, params=params, headers=self._credential.get_headers()), context="写入网盘目录") + + def get_fids(self, file_paths: List[str]) -> List[Dict[str, Any]]: + pending = [self._normalize_dir_path(p) for p in file_paths] + result: List[Dict[str, Any]] = [] + while pending: + batch, pending = pending[:50], pending[50:] + url = "https://drive-pc.quark.cn/1/clouddrive/file/info/path_list" + params = {"pr": "ucpro", "fr": "pc"} + payload = {"file_path": batch, "namespace": "0"} + data = self._drive_api_json(self._post(url, json_data=payload, params=params, headers=self._credential.get_headers()), context="按路径获取文件ID") + if not self._api_success(data): + raise TransferError( + TransferErrorCode.NETWORK_ERROR, + message=f"获取目录ID失败: {data.get('message', data)}", + platform=self.PLATFORM_KEY, + ) + result.extend(data.get("data", [])) + return result + + def move_files(self, fids: List[str], to_pdir_fid: str) -> Dict[str, Any]: + if not fids: + return {"code": 0, "message": "无文件需要移动"} + last: Dict[str, Any] = {"code": 0, "message": "success"} + for offset in range(0, len(fids), 100): + batch = fids[offset:offset + 100] + url = "https://drive-pc.quark.cn/1/clouddrive/file/move" + params = {"uc_param_str": "", "fr": "pc", "pr": "ucpro"} + payload = {"filelist": batch, "to_pdir_fid": to_pdir_fid, "exclude_fids": [], "action_type": 1} + last = self._drive_api_json(self._post(url, json_data=payload, params=params, headers=self._credential.get_headers()), context="移动网盘文件") + if not self._api_success(last): + return last + task_id = last.get("data", {}).get("task_id") + if task_id: + task = self.query_task(task_id) + if not self._api_success(task) or task.get("data", {}).get("status") == -1: + return {"code": 1, "message": task.get("data", {}).get("message", task.get("message", "移动任务失败")), "data": task.get("data", {})} + return {"code": 0, "message": "移动完成", "data": last.get("data", {})} + + def _task_query_params(self, retry_index: int = 0) -> Dict[str, Any]: + now_ms = int(time.time() * 1000) + return { + "pr": "ucpro", + "fr": "pc", + "uc_param_str": "", + "retry_index": retry_index, + "__dt": 300, + "__t": now_ms, + } + + def delete_files(self, fids: List[str]) -> Dict[str, Any]: + if not fids: + return {"code": 0, "status": 200} + if self.delete(fids): + return {"code": 0, "status": 200} + return {"code": 1, "status": 500, "message": "删除文件失败"} + + def query_task(self, task_id: str) -> Dict[str, Any]: + url = "https://drive-pc.quark.cn/1/clouddrive/task" + try: + data = self._poll_task(url, task_id, query_params=self._task_query_params) + return {"code": 0, "status": 200, "data": data} + except TransferError as exc: + return {"code": 1, "status": 500, "message": str(exc), "data": {"status": -1}} + + def recycle_list(self, page: int = 1, size: int = 30) -> List[Dict[str, Any]]: + url = "https://drive-pc.quark.cn/1/clouddrive/file/recycle/list" + params = {"_page": page, "_size": size, "pr": "ucpro", "fr": "pc", "uc_param_str": ""} + data = self._drive_api_json(self._get(url, params=params, headers=self._credential.get_headers()), context="列出回收站") + return data.get("data", {}).get("list", []) + + def recycle_remove(self, record_list: List[Dict[str, Any]]) -> Dict[str, Any]: + url = "https://drive-pc.quark.cn/1/clouddrive/file/recycle/remove" + params = {"uc_param_str": "", "fr": "pc", "pr": "ucpro"} + payload = {"select_mode": 2, "record_list": record_list} + return self._drive_api_json(self._post(url, json_data=payload, params=params, headers=self._credential.get_headers()), context="写入网盘目录") + + def cleanup_recycle(self, fids: List[str]) -> Dict[str, Any]: + target_fids = {str(fid) for fid in fids if fid} + if not target_fids: + return {"code": 0, "message": "无回收站记录需要清理", "data": {"removed": 0}} + records = [ + item for item in self.recycle_list() + if str(item.get("fid") or item.get("file_id") or "") in target_fids + ] + if not records: + return {"code": 0, "message": "未找到匹配的回收站记录", "data": {"removed": 0}} + result = self.recycle_remove(records) + if self._api_success(result): + result.setdefault("data", {})["removed"] = len(records) + return result + # ─── get_files / delete ──────────────────────────────────── def get_files(self, parent_fid: str = "0") -> List[FileInfo]: diff --git a/cloudsearch_transfer/adapter/uc/__init__.py b/cloudsearch_transfer/adapter/uc/__init__.py index 3e87e09..678059f 100644 --- a/cloudsearch_transfer/adapter/uc/__init__.py +++ b/cloudsearch_transfer/adapter/uc/__init__.py @@ -56,19 +56,21 @@ class UcAdapter(BaseCloudDriveAdapter): r"drive\.uc\.cn/s/(\w+)", ] + + capabilities: Dict[str, bool] = { + "ensure_dir": True, + "save_files": True, + "poll_task": True, + "rename": True, + "move_files": True, + "delete_files": True, + "cleanup_recycle": False, + "share_staging_folder": True, + } + def __init__(self, config: PlatformConfig, transfer_config: TransferConfig) -> None: - """初始化 UC 适配器。 - - Args: - config: 平台配置(含 Cookie 等)。 - transfer_config: 全局转存配置(超时、重试、轮询参数等)。 - """ - super().__init__(config, transfer_config) - - # 初始化三个子模块 - self._credential: UcCredentialManager = UcCredentialManager( - cookie=config.cookie - ) + """初始化适配器。""" + self._credential: UcCredentialManager = UcCredentialManager(cookie=config.cookie) self._transfer_engine: UcTransfer = UcTransfer( credential=self._credential, timeout=transfer_config.request_timeout, @@ -79,6 +81,7 @@ class UcAdapter(BaseCloudDriveAdapter): credential=self._credential, timeout=transfer_config.request_timeout, ) + super().__init__(config, transfer_config) # ═══════════════════════════════════════════════════════════════ # 公开接口实现 @@ -93,21 +96,8 @@ class UcAdapter(BaseCloudDriveAdapter): def transfer(self, share_url: str, save_dir: str = "", share_password: str = "") -> TransferResult: - """执行转存的核心逻辑(覆盖基类实现 UC 专用流程)。 - - 通过 UcTransfer 引擎执行完整的 7 步流程。 - - Args: - share_url: UC 分享链接。 - save_dir: 目标目录,空则使用配置的默认目录。 - share_password: 新分享的密码。 - - Returns: - TransferResult 包含转存结果。 - """ start: float = time.time() - # 凭证检查 if not self._credential.validate(): raise TransferError( TransferErrorCode.NOT_LOGIN, @@ -115,18 +105,36 @@ class UcAdapter(BaseCloudDriveAdapter): platform=self.PLATFORM_KEY, ) - # 目标目录:默认根目录 "0" - target_dir: str = save_dir or self.config.save_dir or "0" - - # 分享密码 + requested_dir: str = save_dir or self.config.save_dir or "/" + target_dir: str = requested_dir if requested_dir and not str(requested_dir).startswith("/") else self.ensure_dir(requested_dir) + staging_dir: str = self.get_or_create_share_folder() or target_dir pwd: str = share_password or self.config.share_password or "" try: - result: Dict[str, Any] = self._transfer_engine.transfer( - share_url=share_url, - save_dir=target_dir, - share_password=pwd, + pwd_id, passcode = self._parse_share_url(share_url) + stoken: str = self._transfer_engine._get_stoken(pwd_id, passcode) + detail: Dict[str, Any] = self._transfer_engine._get_detail(pwd_id, stoken) + task_id: str = self._transfer_engine._init_save( + pwd_id, stoken, detail, to_pdir_fid=staging_dir ) + new_fids: List[str] = self._transfer_engine._poll_save_task(task_id) + if not new_fids: + raise RuntimeError("转存完成但未获取到文件ID") + + if staging_dir != target_dir: + move_result = self.move_files(new_fids, target_dir) + if not self._api_success(move_result): + raise RuntimeError(f"移动到目标目录失败: {move_result.get('message', move_result)}") + + if self.transfer_config.ad_filter_enabled and new_fids: + new_fids = self._filter_ads(new_fids) + if not new_fids: + raise RuntimeError("广告过滤后无可分享文件") + + title: str = detail.get("title", "分享") + share_task_id: str = self._transfer_engine._init_share(new_fids, title) + share_id: str = self._transfer_engine._poll_share_task(share_task_id) + share_url_new, passcode_new = self._transfer_engine._set_password(share_id, pwd) except ValueError as exc: raise TransferError( TransferErrorCode.URL_INVALID, @@ -148,24 +156,13 @@ class UcAdapter(BaseCloudDriveAdapter): ) from exc elapsed: int = int((time.time() - start) * 1000) - - # 广告过滤 - new_fids: List[str] = result.get("new_file_ids", []) - if self.transfer_config.ad_filter_enabled and new_fids: - new_fids = self._filter_ads(new_fids) - if not new_fids: - raise TransferError( - TransferErrorCode.RESOURCE_EMPTY, - platform=self.PLATFORM_KEY, - ) - return TransferResult( success=True, platform=self.PLATFORM_KEY, new_file_id=",".join(new_fids), - file_name=result.get("file_name", ""), - share_url=result.get("share_url", ""), - share_password=result.get("passcode", pwd), + file_name=title, + share_url=share_url_new, + share_password=passcode_new, original_url=share_url, elapsed_ms=elapsed, ) @@ -204,8 +201,8 @@ class UcAdapter(BaseCloudDriveAdapter): files=files, ) - except TransferError: - raise + except TransferError as exc: + return VerifyResult(valid=False, platform=self.PLATFORM_KEY, error=exc) except (ValueError, RuntimeError) as exc: return VerifyResult( valid=False, @@ -344,6 +341,154 @@ class UcAdapter(BaseCloudDriveAdapter): ) return file_ids + + # ─── Drive API capability helpers ───────────────────────────── + + @staticmethod + def _normalize_dir_path(dir_path: str) -> str: + path = "/" + str(dir_path or "").strip().strip("/") + return "/" if path == "/" else path + + @staticmethod + def _api_success(payload: Dict[str, Any]) -> bool: + if not isinstance(payload, dict): + return False + code = payload.get("code") + status = payload.get("status") + return code == 0 or status in (0, 200) + + def ensure_dir(self, dir_path: str) -> str: + normalized = self._normalize_dir_path(dir_path) + if normalized == "/": + return "0" + + parts = [part for part in normalized.strip("/").split("/") if part] + prefixes = ["/" + "/".join(parts[:idx]) for idx in range(1, len(parts) + 1)] + existing = { + item.get("file_path"): str(item.get("fid")) + for item in self.get_fids(prefixes) + if item.get("file_path") and item.get("fid") + } + + leaf_fid = existing.get(normalized) + if leaf_fid: + return leaf_fid + + for prefix in prefixes: + if prefix in existing: + continue + result = self.mkdir(prefix) + if self._api_success(result) and result.get("data", {}).get("fid"): + existing[prefix] = str(result["data"]["fid"]) + continue + raise TransferError( + TransferErrorCode.NETWORK_ERROR, + message=f"创建目录失败: {result.get('message', result)}", + platform=self.PLATFORM_KEY, + ) + return existing[normalized] + + def mkdir(self, dir_path: str) -> Dict[str, Any]: + url = "https://pc-api.uc.cn/1/clouddrive/file" + params = {"pr": "UCBrowser", "fr": "pc", "uc_param_str": ""} + payload = { + "pdir_fid": "0", + "file_name": "", + "dir_path": self._normalize_dir_path(dir_path), + "dir_init_lock": False, + } + return self._drive_api_json(self._post(url, json_data=payload, params=params, headers=self._credential.get_headers()), context="写入网盘目录") + + def rename(self, fid: str, file_name: str) -> Dict[str, Any]: + url = "https://pc-api.uc.cn/1/clouddrive/file/rename" + params = {"pr": "UCBrowser", "fr": "pc", "uc_param_str": ""} + payload = {"fid": fid, "file_name": file_name} + return self._drive_api_json(self._post(url, json_data=payload, params=params, headers=self._credential.get_headers()), context="写入网盘目录") + + def get_fids(self, file_paths: List[str]) -> List[Dict[str, Any]]: + pending = [self._normalize_dir_path(p) for p in file_paths] + result: List[Dict[str, Any]] = [] + while pending: + batch, pending = pending[:50], pending[50:] + url = "https://pc-api.uc.cn/1/clouddrive/file/info/path_list" + params = {"pr": "UCBrowser", "fr": "pc"} + payload = {"file_path": batch, "namespace": "0"} + data = self._drive_api_json(self._post(url, json_data=payload, params=params, headers=self._credential.get_headers()), context="按路径获取文件ID") + if not self._api_success(data): + raise TransferError( + TransferErrorCode.NETWORK_ERROR, + message=f"获取目录ID失败: {data.get('message', data)}", + platform=self.PLATFORM_KEY, + ) + result.extend(data.get("data", [])) + return result + + def move_files(self, fids: List[str], to_pdir_fid: str) -> Dict[str, Any]: + if not fids: + return {"code": 0, "message": "无文件需要移动"} + last: Dict[str, Any] = {"code": 0, "message": "success"} + for offset in range(0, len(fids), 100): + batch = fids[offset:offset + 100] + url = "https://pc-api.uc.cn/1/clouddrive/file/move" + params = {"uc_param_str": "", "fr": "pc", "pr": "UCBrowser"} + payload = {"filelist": batch, "to_pdir_fid": to_pdir_fid, "exclude_fids": [], "action_type": 1} + last = self._drive_api_json(self._post(url, json_data=payload, params=params, headers=self._credential.get_headers()), context="移动网盘文件") + if not self._api_success(last): + return last + task_id = last.get("data", {}).get("task_id") + if task_id: + task = self.query_task(task_id) + if not self._api_success(task) or task.get("data", {}).get("status") == -1: + return {"code": 1, "message": task.get("data", {}).get("message", task.get("message", "移动任务失败")), "data": task.get("data", {})} + return {"code": 0, "message": "移动完成", "data": last.get("data", {})} + + def _task_query_params(self, retry_index: int = 0) -> Dict[str, Any]: + now_ms = int(time.time() * 1000) + return { + "pr": "UCBrowser", + "fr": "pc", + "uc_param_str": "", + "retry_index": retry_index, + "__dt": 300, + "__t": now_ms, + } + + def delete_files(self, fids: List[str]) -> Dict[str, Any]: + if not fids: + return {"code": 0, "status": 200} + if self.delete(fids): + return {"code": 0, "status": 200} + return {"code": 1, "status": 500, "message": "删除文件失败"} + + def query_task(self, task_id: str) -> Dict[str, Any]: + url = "https://pc-api.uc.cn/1/clouddrive/task" + try: + data = self._poll_task(url, task_id, query_params=self._task_query_params) + return {"code": 0, "status": 200, "data": data} + except TransferError as exc: + return {"code": 1, "status": 500, "message": str(exc), "data": {"status": -1}} + + def get_or_create_share_folder(self) -> Optional[str]: + if getattr(self, "_share_folder_fid", None): + return self._share_folder_fid + root = self.ls_dir("0") + if self._api_success(root): + for item in root.get("data", {}).get("list", []): + if item.get("file_name") == "来自:分享" and item.get("dir"): + self._share_folder_fid = str(item["fid"]) + return self._share_folder_fid + result = self.mkdir("/来自:分享") + if self._api_success(result) and result.get("data", {}).get("fid"): + self._share_folder_fid = str(result["data"]["fid"]) + return self._share_folder_fid + return None + + def ls_dir(self, pdir_fid: str) -> Dict[str, Any]: + url = "https://pc-api.uc.cn/1/clouddrive/file/sort" + params = {"pr": "UCBrowser", "fr": "pc", "pdir_fid": pdir_fid or "0", "_page": 1, "_size": 50, "_fetch_total": 1, "_fetch_sub_dirs": 0, "_sort": "file_type:asc,updated_at:desc"} + return self._drive_api_json(self._get(url, params=params, headers=self._credential.get_headers()), context="列出网盘目录") + + # ─── get_files / delete ──────────────────────────────────── def get_files(self, parent_fid: str = "0") -> List[FileInfo]: diff --git a/cloudsearch_transfer/tests/test_quark_uc_capabilities_unittest.py b/cloudsearch_transfer/tests/test_quark_uc_capabilities_unittest.py new file mode 100644 index 0000000..a4b99ea --- /dev/null +++ b/cloudsearch_transfer/tests/test_quark_uc_capabilities_unittest.py @@ -0,0 +1,308 @@ + +import unittest +import requests +from unittest.mock import patch + +from cloudsearch_transfer.adapter.quark import QuarkAdapter +from cloudsearch_transfer.adapter.uc import UcAdapter +from cloudsearch_transfer.config import PlatformConfig, TransferConfig +from cloudsearch_transfer.adapter.base import BaseCloudDriveAdapter +from cloudsearch_transfer.errors import TransferError + + +class DummyResponse: + def __init__(self, payload): + self._payload = payload + def json(self): + return self._payload + def raise_for_status(self): + return None + + +def make_adapter(adapter_cls): + return adapter_cls( + PlatformConfig(enabled=True, cookie="k=" + "x" * 80, account_name="test"), + TransferConfig(request_timeout=1, max_retries=0, ad_filter_enabled=False), + ) + + +class QuarkUcCapabilityTests(unittest.TestCase): + def test_quark_declares_p0_drive_api_capabilities(self): + adapter = make_adapter(QuarkAdapter) + self.assertTrue(adapter.capabilities["ensure_dir"]) + self.assertTrue(adapter.capabilities["save_files"]) + self.assertTrue(adapter.capabilities["poll_task"]) + self.assertTrue(adapter.capabilities["rename"]) + self.assertTrue(adapter.capabilities["move_files"]) + self.assertTrue(adapter.capabilities["delete_files"]) + self.assertTrue(adapter.capabilities["cleanup_recycle"]) + + def test_uc_declares_p0_drive_api_capabilities_and_share_staging(self): + adapter = make_adapter(UcAdapter) + self.assertTrue(adapter.capabilities["ensure_dir"]) + self.assertTrue(adapter.capabilities["save_files"]) + self.assertTrue(adapter.capabilities["poll_task"]) + self.assertTrue(adapter.capabilities["rename"]) + self.assertTrue(adapter.capabilities["move_files"]) + self.assertTrue(adapter.capabilities["delete_files"]) + self.assertTrue(adapter.capabilities["share_staging_folder"]) + + def test_quark_ensure_dir_uses_path_lookup_before_mkdir(self): + adapter = make_adapter(QuarkAdapter) + mkdir_calls = [] + adapter.get_fids = lambda paths: [{"file_path": "/Media", "fid": "fid-media"}] + adapter.mkdir = lambda path: mkdir_calls.append(path) or {"code": 0, "data": {"fid": "new"}} + self.assertEqual(adapter.ensure_dir("/Media"), "fid-media") + self.assertEqual(mkdir_calls, []) + + def test_quark_ensure_dir_creates_missing_path(self): + adapter = make_adapter(QuarkAdapter) + adapter.get_fids = lambda paths: [] + adapter.mkdir = lambda path: {"code": 0, "data": {"fid": "created-fid"}} + self.assertEqual(adapter.ensure_dir("Media/Shows"), "created-fid") + + def test_uc_share_folder_created_when_missing(self): + adapter = make_adapter(UcAdapter) + adapter.ls_dir = lambda parent: {"code": 0, "data": {"list": []}} + adapter.mkdir = lambda path: {"code": 0, "data": {"fid": "share-folder-fid"}} + self.assertEqual(adapter.get_or_create_share_folder(), "share-folder-fid") + + def test_uc_move_files_polls_task_until_complete(self): + adapter = make_adapter(UcAdapter) + calls = [] + def fake_post(url, json_data=None, params=None, headers=None): + calls.append((url, json_data)) + return DummyResponse({"code": 0, "status": 200, "data": {"task_id": "task-1"}}) + adapter._post = fake_post + adapter.query_task = lambda task_id: {"code": 0, "status": 200, "data": {"status": 2}} + result = adapter.move_files(["fid-1", "fid-2"], "target-fid") + self.assertEqual(result["code"], 0) + self.assertEqual(calls[0][1]["filelist"], ["fid-1", "fid-2"]) + self.assertEqual(calls[0][1]["to_pdir_fid"], "target-fid") + + def test_quark_save_files_resolves_path_to_fid(self): + adapter = make_adapter(QuarkAdapter) + adapter.ensure_dir = lambda path: "target-fid" + adapter._transfer_engine._get_stoken = lambda pwd_id: "stoken" + captured = {} + adapter._transfer_engine._init_save = lambda pwd_id, stoken, detail, to_pdir_fid: captured.setdefault("to_pdir_fid", to_pdir_fid) or "task-1" + adapter._transfer_engine._poll_save_task = lambda task_id: ["new-fid"] + self.assertEqual(adapter._save_files("pwd", {"fid": "src"}, "/Media"), ["new-fid"]) + self.assertEqual(captured["to_pdir_fid"], "target-fid") + + def test_uc_transfer_resolves_path_to_fid(self): + adapter = make_adapter(UcAdapter) + adapter._credential.validate = lambda: True + adapter.ensure_dir = lambda path: "target-fid" + adapter.get_or_create_share_folder = lambda: "target-fid" + adapter._parse_share_url = lambda url: ("pwd", "pass") + adapter._transfer_engine._get_stoken = lambda pwd_id, passcode="": "stoken" + adapter._transfer_engine._get_detail = lambda pwd_id, stoken: {"title": "demo", "fid": "src"} + captured = {} + def fake_init_save(pwd_id, stoken, detail, to_pdir_fid): + captured["save_dir"] = to_pdir_fid + return "save-task" + adapter._transfer_engine._init_save = fake_init_save + adapter._transfer_engine._poll_save_task = lambda task_id: ["new-fid"] + adapter._transfer_engine._init_share = lambda fids, title: "share-task" + adapter._transfer_engine._poll_share_task = lambda task_id: "share-id" + adapter._transfer_engine._set_password = lambda share_id, password: ("https://drive.uc.cn/s/new", password) + result = adapter.transfer("https://drive.uc.cn/s/abcdef", save_dir="/Media", share_password="pw") + self.assertTrue(result.success) + self.assertEqual(captured["save_dir"], "target-fid") + def test_quark_transfer_resolves_path_to_fid(self): + adapter = make_adapter(QuarkAdapter) + adapter.ensure_dir = lambda path: "target-fid" + adapter._credential.validate = lambda: True + captured = {} + def fake_transfer(share_url, save_dir, share_password): + captured["save_dir"] = save_dir + return { + "new_file_ids": ["new-fid"], + "file_name": "demo", + "share_url": "https://pan.quark.cn/s/new", + "passcode": share_password, + } + adapter._transfer_engine.transfer = fake_transfer + result = adapter._transfer("https://pan.quark.cn/s/abcdef", save_dir="/Media", share_password="pw") + self.assertTrue(result.success) + self.assertEqual(captured["save_dir"], "target-fid") + + def test_ensure_dir_creates_nested_paths_progressively(self): + adapter = make_adapter(QuarkAdapter) + adapter.get_fids = lambda paths: [] + created = [] + def fake_mkdir(path): + created.append(path) + return {"code": 0, "data": {"fid": "fid-" + path.rsplit("/", 1)[-1]}} + adapter.mkdir = fake_mkdir + self.assertEqual(adapter.ensure_dir("/A/B"), "fid-B") + self.assertEqual(created, ["/A", "/A/B"]) + + def test_uc_transfer_uses_staging_folder_then_moves_to_target(self): + adapter = make_adapter(UcAdapter) + adapter._credential.validate = lambda: True + adapter.ensure_dir = lambda path: "target-fid" + adapter.get_or_create_share_folder = lambda: "staging-fid" + adapter._parse_share_url = lambda url: ("pwd", "pass") + adapter._transfer_engine._get_stoken = lambda pwd_id, passcode="": "stoken" + adapter._transfer_engine._get_detail = lambda pwd_id, stoken: {"title": "demo", "fid": "src"} + captured = {} + def fake_init_save(pwd_id, stoken, detail, to_pdir_fid): + captured["save_to"] = to_pdir_fid + return "save-task" + adapter._transfer_engine._init_save = fake_init_save + adapter._transfer_engine._poll_save_task = lambda task_id: ["new-fid"] + def fake_move(fids, to_pdir_fid): + captured["move_to"] = to_pdir_fid + return {"code": 0, "status": 200} + adapter.move_files = fake_move + adapter._transfer_engine._init_share = lambda fids, title: "share-task" + adapter._transfer_engine._poll_share_task = lambda task_id: "share-id" + adapter._transfer_engine._set_password = lambda share_id, password: ("https://drive.uc.cn/s/new", password) + result = adapter.transfer("https://drive.uc.cn/s/abcdef", save_dir="/Media", share_password="pw") + self.assertTrue(result.success) + self.assertEqual(captured["save_to"], "staging-fid") + self.assertEqual(captured["move_to"], "target-fid") + + def test_api_success_requires_explicit_success_code_or_status(self): + adapter = make_adapter(QuarkAdapter) + self.assertFalse(adapter._api_success({})) + self.assertFalse(adapter._api_success({"message": "bad"})) + self.assertTrue(adapter._api_success({"code": 0})) + self.assertTrue(adapter._api_success({"status": 200})) + + def test_query_task_passes_retry_index_and_timestamp_params(self): + adapter = make_adapter(QuarkAdapter) + captured_params = [] + attempts = {"count": 0} + def fake_get(url, params=None, retry=None): + captured_params.append(dict(params)) + attempts["count"] += 1 + if attempts["count"] == 1: + return DummyResponse({"data": {"status": 1}}) + return DummyResponse({"data": {"status": 2, "task_id": "task-1"}}) + adapter._get = fake_get + result = adapter.query_task("task-1") + self.assertEqual(result["code"], 0) + self.assertEqual(captured_params[0]["retry_index"], 0) + self.assertEqual(captured_params[1]["retry_index"], 1) + self.assertIn("__dt", captured_params[0]) + self.assertIn("__t", captured_params[0]) + + def test_quark_cleanup_recycle_only_removes_matching_fids(self): + adapter = make_adapter(QuarkAdapter) + adapter.recycle_list = lambda: [ + {"fid": "keep", "record_id": "r-keep"}, + {"fid": "target", "record_id": "r-target"}, + ] + captured = {} + def fake_recycle_remove(records): + captured["records"] = records + return {"code": 0} + adapter.recycle_remove = fake_recycle_remove + result = adapter.cleanup_recycle(["target"]) + self.assertEqual(result["code"], 0) + self.assertEqual(captured["records"], [{"fid": "target", "record_id": "r-target"}]) + + def test_drive_api_json_raises_transfer_error_on_http_error(self): + adapter = make_adapter(QuarkAdapter) + class ErrorResponse(DummyResponse): + status_code = 429 + text = "too many requests" + def raise_for_status(self): + raise requests.HTTPError("429 Too Many Requests") + with self.assertRaises(Exception) as ctx: + adapter._drive_api_json(ErrorResponse({"code": 0}), context="限流测试") + self.assertIn("限流测试", str(ctx.exception)) + + def test_drive_api_json_raises_transfer_error_on_invalid_json(self): + adapter = make_adapter(QuarkAdapter) + class InvalidJsonResponse(DummyResponse): + status_code = 200 + text = "bad gateway" + def json(self): + raise ValueError("not json") + with self.assertRaises(Exception) as ctx: + adapter._drive_api_json(InvalidJsonResponse({}), context="JSON测试") + self.assertIn("JSON测试", str(ctx.exception)) + + + def test_quark_delete_files_capability_calls_drive_delete(self): + adapter = make_adapter(QuarkAdapter) + calls = [] + adapter.delete = lambda fids: calls.append(list(fids)) or True + self.assertEqual(adapter.delete_files(["fid-1", "fid-2"]), {"code": 0, "status": 200}) + self.assertEqual(calls, [["fid-1", "fid-2"]]) + + def test_uc_delete_files_capability_calls_drive_delete(self): + adapter = make_adapter(UcAdapter) + calls = [] + adapter.delete = lambda fids: calls.append(list(fids)) or True + self.assertEqual(adapter.delete_files(["fid-1"]), {"code": 0, "status": 200}) + self.assertEqual(calls, [["fid-1"]]) + + def test_uc_transfer_filters_ads_after_staging_save(self): + adapter = make_adapter(UcAdapter) + adapter.ensure_dir = lambda path: "target-dir" + adapter.get_or_create_share_folder = lambda: "staging-dir" + adapter._parse_share_url = lambda url: ("pwd-id", "") + adapter._filter_ads = lambda fids: [fid for fid in fids if fid != "ad-fid"] + adapter.transfer_config.ad_filter_enabled = True + adapter._transfer_engine._get_stoken = lambda pwd_id, passcode: "stoken" + adapter._transfer_engine._get_detail = lambda pwd_id, stoken: {"title": "Title"} + adapter._transfer_engine._init_save = lambda pwd_id, stoken, detail, to_pdir_fid: "save-task" + adapter._transfer_engine._poll_save_task = lambda task_id: ["keep-fid", "ad-fid"] + adapter.move_files = lambda fids, target: {"code": 0} + adapter._transfer_engine._init_share = lambda fids, title: "share-task" if fids == ["keep-fid"] else (_ for _ in ()).throw(AssertionError(f"unfiltered fids: {fids}")) + adapter._transfer_engine._poll_share_task = lambda task_id: "share-id" + adapter._transfer_engine._set_password = lambda share_id, pwd: ("https://share", pwd) + + result = adapter.transfer("https://drive.uc.cn/s/abc") + + self.assertEqual(result.new_file_id, "keep-fid") + + def test_uc_ls_dir_wraps_http_and_json_errors(self): + adapter = make_adapter(UcAdapter) + class BadResponse: + text = "bad gateway" + def raise_for_status(self): + raise requests.HTTPError("HTTP 502") + def json(self): + raise AssertionError("json() should not be called after HTTP error") + adapter._get = lambda *args, **kwargs: BadResponse() + with self.assertRaises(TransferError): + adapter.ls_dir("0") + + def test_base_declares_optional_drive_api_methods(self): + required = [ + "ensure_dir", "get_fids", "mkdir", "rename", "move_files", + "delete_files", "query_task", "cleanup_recycle", + ] + for name in required: + self.assertTrue(hasattr(BaseCloudDriveAdapter, name), name) + + def test_base_optional_drive_api_methods_raise_transfer_error(self): + class MinimalAdapter(BaseCloudDriveAdapter): + PLATFORM_KEY = "minimal" + def verify(self, share_url): + raise NotImplementedError + def _get_share_detail(self, pwd_id, passcode=""): + raise NotImplementedError + def _save_files(self, pwd_id, detail, save_dir): + raise NotImplementedError + def _create_share(self, file_ids, title, password=""): + raise NotImplementedError + def get_files(self, parent_fid="0"): + raise NotImplementedError + def delete(self, file_ids): + raise NotImplementedError + adapter = MinimalAdapter(PlatformConfig(), TransferConfig()) + with self.assertRaises(TransferError): + adapter.ensure_dir("/Media") + with self.assertRaises(TransferError): + adapter.move_files(["fid"], "target") + + +if __name__ == "__main__": + unittest.main() diff --git a/docs/deployment/test-server.md b/docs/deployment/test-server.md new file mode 100644 index 0000000..095f368 --- /dev/null +++ b/docs/deployment/test-server.md @@ -0,0 +1,19 @@ +# CloudSearch 测试服部署记录 + +- 服务器:泽御云香港测试服,SSH 入口 root@82.158.228.152:18924。 +- 当前运行来源:`/root/cloudsearch_deploy`。 +- 主应用目录:`/root/cloudsearch_deploy/source_clean`。 +- 网盘能力层:`/root/cloudsearch_deploy/cloudsearch_transfer`。 +- 启动方式:Docker Compose,主容器挂载 `/root/cloudsearch_deploy` 相关目录。 +- 集成边界:只集成 cloud-auto-save 中可复用的网盘 API 调用能力;不集成签到、收益事件、收益上报,也不长期部署独立 cloud-auto-save 服务。 +- Video Parser 当前未部署,属于测试环境预期缺口,不作为主服务故障处理。 +- 本轮隔离点:`backup-before-drive-api-20260522-152138`;开发分支:`feature/cloudsearch-drive-api-20260522-152138`。 + +## 验证命令 + +```bash +python3 -m unittest discover -s cloudsearch_transfer/tests -p 'test_*.py' -v +python3 /tmp/verify_compile.py +docker compose ps +curl -fsS http://127.0.0.1:3000/health +```