import json import logging import os import subprocess import time import uuid from typing import Any, Dict, List, Literal, Optional, TypedDict import requests from DrissionPage import Chromium from DrissionPage._pages.chromium_tab import ChromiumTab from DrissionPage.common import By try: import winreg except ImportError: winreg = None logger = logging.getLogger(__name__) STATUS_OK = "0" STATUS_LOGIN_FAILED = "-10003" DEFAULT_SOCKET_PORT = 19890 CLIENT_API_TIMEOUT = 120 PORT_CHECK_TIMEOUT = 2 UPDATE_CORE_RETRY_DELAY = 2 CLIENT_RESTART_DELAY = 5 CLIENT_START_RETRIES = 3 CLIENT_READY_TIMEOUT = 10 CLIENT_READY_INTERVAL = 0.5 CLIENT_POST_START_DELAY = 5 PROCESS_KILL_DELAY = 3 DOC_LOAD_TIMEOUT = 30 COUNTRY_INITIAL_LOAD_TIMEOUT = 120 COUNTRY_LOOKUP_TIMEOUT = 20 COUNTRY_CLICK_TIMEOUT = 10 COUNTRY_DROPDOWN_DELAY = 1 IP_CHECK_TIMEOUT = 60 LOGIN_ATTEMPTS = 4 PASSWORD_INPUT_TIMEOUT = 5 PASSWORD_SUBMIT_TIMEOUT = 5 OTP_SEND_TIMEOUT = 5 OTP_INPUT_LOOKUP_ATTEMPTS = 5 OTP_INPUT_TIMEOUT = 30 OTP_INPUT_DISPLAY_TIMEOUT = 10 OTP_CODE_WAIT_SECONDS = 30 OTP_SUBMIT_TIMEOUT = 10 OTP_RESULT_LOAD_TIMEOUT = 20 OTP_ERROR_TIMEOUT = 10 OTP_SEND_DELAY = 1 LOGIN_FINAL_SUBMIT_TIMEOUT = 10 NEED_LOGIN_DELAY = 3 NEED_LOGIN_TIMEOUT = 5 COUNTRY_LABEL_XPATH = 'xpath://div[@class="dropdown-account-switcher-header-label"]/span[last()]' COUNTRY_DROPDOWN_XPATH = 'xpath://div[@class="dropdown-account-switcher-header-label"]' COUNTRY_LIST_ITEM_XPATH = 'xpath://div[@class="dropdown-account-switcher-list-item"]' NEED_LOGIN_XPATH = 'xpath://h1[@class="a-spacing-small"]|//span[contains(text(),"登录")]' PASSWORD_INPUT_XPATH = 'xpath://input[@type="password"]' PASSWORD_SUBMIT_XPATH = 'xpath://input[@id="signInSubmit"]' OTP_SEND_XPATH = 'xpath://span[@id="auth-send-code" and contains(string(.),"发送一次性密码")]' OTP_INPUT_XPATH = 'xpath://input[@name="otpCode"]' OTP_SUBMIT_XPATH = 'xpath://input[@id="auth-signin-button"]' OTP_ERROR_XPATH = 'xpath://div[@id="auth-error-message-box"]' class UserInfo(TypedDict): username: str password: str company: str def kill_process(version: Literal["v5", "v6"]): """结束指定版本的紫鸟客户端进程.""" logger.info("准备杀紫鸟客户端进程,version=%s", version) driver = ZiniaoDriver({}) driver.kill_process(version) class ZiniaoDriver: """封装紫鸟客户端启动,店铺浏览器生命周期和客户端 API 调用.""" def __init__(self, user_info: UserInfo, socket_port: int = DEFAULT_SOCKET_PORT): """初始化紫鸟浏览器驱动. Args: user_info: 紫鸟账号信息,包含 company, username, password. socket_port: 客户端 HTTP 通信端口. """ self.user_info = user_info self.socket_port = socket_port self.client_path = None self.browser = None self.tab: Optional[ChromiumTab] = None self.store_id = None @property def client_url(self) -> str: return f"http://127.0.0.1:{self.socket_port}" def _build_payload(self, action: str, extra: Optional[Dict[str, Any]] = None) -> Dict[str, Any]: payload = { "action": action, "requestId": str(uuid.uuid4()), } if extra: payload.update(extra) payload.update(self.user_info) return payload def _post_client(self, payload: Dict[str, Any]) -> Dict[str, Any]: response = requests.post( self.client_url, json.dumps(payload).encode("utf-8"), timeout=CLIENT_API_TIMEOUT, ) return response.json() @staticmethod def _status_code(result: Optional[Dict[str, Any]]) -> Optional[str]: if result is None: return None return str(result.get("statusCode")) @staticmethod def _result_text(result: Dict[str, Any]) -> str: return json.dumps(result, ensure_ascii=False) @staticmethod def _find_store_oauth(shop_list: Optional[List[Dict[str, Any]]], shop_name: str): for shop in shop_list or []: if shop.get("browserName") == shop_name: return shop.get("browserOauth") return None def _is_client_ready(self) -> bool: try: requests.get(self.client_url, timeout=PORT_CHECK_TIMEOUT) return True except (requests.exceptions.ConnectionError, requests.exceptions.Timeout): return False def _wait_until_client_ready(self, timeout: float) -> bool: start_check_time = time.time() while time.time() - start_check_time < timeout: if self._is_client_ready(): return True time.sleep(CLIENT_READY_INTERVAL) return False def get_zinaio_exe(self, protocol_name: str = "superbrowser"): """从 Windows 注册表读取紫鸟客户端可执行文件路径. Args: protocol_name: 注册表协议名称. Returns: 紫鸟浏览器可执行文件路径. 未找到时返回 None. """ if winreg is None: logger.error("当前系统不支持 winreg,无法从注册表获取紫鸟客户端路径") return None key_path = rf"SOFTWARE\Classes\{protocol_name}\shell\open\command" for root_key, root_name in ( (winreg.HKEY_CURRENT_USER, "HKEY_CURRENT_USER"), (winreg.HKEY_LOCAL_MACHINE, "HKEY_LOCAL_MACHINE"), ): try: logger.info("正在从注册表读取紫鸟客户端路径:%s\\%s", root_name, key_path) key = winreg.OpenKey(root_key, key_path) command, _ = winreg.QueryValueEx(key, "") winreg.CloseKey(key) if isinstance(command, str): sub = "ziniao.exe" exe_path = command[0 : command.find(sub) + len(sub) + 1] else: exe_path = command[0] logger.info("已获取紫鸟客户端路径:%s", exe_path) return exe_path except FileNotFoundError: logger.warning("注册表中未找到紫鸟客户端路径:%s\\%s", root_name, key_path) logger.error("未能获取紫鸟客户端路径,protocol_name=%s", protocol_name) return None def update_core(self): """更新紫鸟浏览器内核. 打开店铺前调用,要求客户端版本 5.285.7 以上. 接口可能因 HTTP 超时 返回未完成状态,因此会循环调用直到成功或判定客户端不支持. """ payload = self._build_payload("updateCore") logger.info("开始更新紫鸟浏览器内核") while True: result = self._post_client(payload) logger.info("更新内核返回:%s", result) if self._handle_update_core_result(result): return if result is None: continue logger.info("等待更新内核完成:%s", self._result_text(result)) time.sleep(UPDATE_CORE_RETRY_DELAY) def _handle_update_core_result(self, result: Optional[Dict[str, Any]]) -> bool: if result is None: logger.info("等待客户端启动后继续更新内核") time.sleep(UPDATE_CORE_RETRY_DELAY) return False status_code = self._status_code(result) if status_code is None or status_code == STATUS_LOGIN_FAILED: logger.error("当前紫鸟客户端版本不支持更新内核接口,请升级客户端") return True if status_code == STATUS_OK: logger.info("紫鸟浏览器内核更新完成") return True return False def kill_process(self, version: Literal["v5", "v6"]): """结束指定版本的紫鸟客户端进程. Args: version: 客户端主版本. """ if version == "v5": process_name = "SuperBrowser.exe" logger.info("准备结束紫鸟 v5 starter.exe") os.system("taskkill /f /t /im starter.exe") else: process_name = "ziniao.exe" logger.info("准备结束紫鸟客户端进程:%s", process_name) os.system("taskkill /f /t /im " + process_name) time.sleep(PROCESS_KILL_DELAY) def get_browser_list(self) -> Optional[list[Dict[str, Any]]]: """获取紫鸟浏览器店铺列表. Returns: 浏览器店铺列表.登录失败或请求失败时返回 None. """ payload = self._build_payload("getBrowserList") request_id = payload["requestId"] logger.info("开始获取紫鸟浏览器列表,requestId=%s", request_id) result = self._post_client(payload) status_code = self._status_code(result) if status_code == STATUS_OK: browser_list = result.get("browserList") logger.info("获取紫鸟浏览器列表成功,数量=%s", len(browser_list or [])) return browser_list if status_code == STATUS_LOGIN_FAILED: logger.error("获取紫鸟浏览器列表登录失败:%s", self._result_text(result)) return None logger.error("获取紫鸟浏览器列表失败:%s", self._result_text(result)) return None def open_store( self, store_info, isWebDriverReadOnlyMode=0, isprivacy=0, isHeadless=0, cookieTypeSave=0, jsInfo="", ): """启动紫鸟店铺浏览器. Args: store_info: 店铺 browserId 或 browserOauth. isWebDriverReadOnlyMode: WebDriver 只读模式开关. isprivacy: 隐私模式开关. isHeadless: 无头模式开关. cookieTypeSave: Cookie 保存类型. jsInfo: 注入的 JS 信息. Returns: 紫鸟客户端返回的启动结果. """ payload = self._build_start_browser_payload( store_info=store_info, is_web_driver_read_only_mode=isWebDriverReadOnlyMode, is_privacy=isprivacy, is_headless=isHeadless, cookie_type_save=cookieTypeSave, js_info=jsInfo, ) request_id = payload["requestId"] logger.info("开始打开紫鸟店铺,store_info=%s,requestId=%s", store_info, request_id) result = self._post_client(payload) status_code = self._status_code(result) if status_code == STATUS_OK: logger.info("紫鸟店铺打开成功,debuggingPort=%s", result.get("debuggingPort")) return result if status_code == STATUS_LOGIN_FAILED: logger.error("打开紫鸟店铺登录失败:%s", self._result_text(result)) raise RuntimeError(f"[open_store]登录失败 {self._result_text(result)}") logger.error("打开紫鸟店铺失败:%s", self._result_text(result)) raise RuntimeError(f"[open_store]失败 {self._result_text(result)} ") def _build_start_browser_payload( self, *, store_info, is_web_driver_read_only_mode=0, is_privacy=0, is_headless=0, cookie_type_save=0, js_info="", ) -> Dict[str, Any]: payload = self._build_payload( "startBrowser", { "isWaitPluginUpdate": 0, "isHeadless": is_headless, "isWebDriverReadOnlyMode": is_web_driver_read_only_mode, "cookieTypeLoad": 0, "cookieTypeSave": cookie_type_save, "runMode": "1", "isLoadUserPlugin": True, "pluginIdType": 1, "privacyMode": is_privacy, }, ) if store_info.isdigit(): payload["browserId"] = store_info else: payload["browserOauth"] = store_info if len(str(js_info)) > 2: payload["injectJsInfo"] = json.dumps(js_info) return payload def get_browser(self, port) -> Chromium: """连接 DrissionPage 浏览器实例. Args: port: Chromium 调试端口. Returns: DrissionPage Chromium 实例. """ logger.info("开始连接 DrissionPage 浏览器,port=%s", port) browser = Chromium(port) logger.info("DrissionPage 浏览器连接完成,port=%s", port) return browser def start_client(self): """启动紫鸟客户端并更新浏览器内核.""" if self._is_client_ready(): logger.info("端口 %s 已启动,跳过启动紫鸟客户端", self.socket_port) return logger.info("端口 %s 未启动,开始启动紫鸟客户端", self.socket_port) self.kill_process("v6") time.sleep(CLIENT_RESTART_DELAY) client_path = self.get_zinaio_exe("superbrowserv6") if not client_path: raise RuntimeError("未找到紫鸟客户端路径,无法启动客户端") self.client_path = client_path.strip('"') cmd = [ self.client_path, "--run_type=web_driver", "--ipc_type=http", "--port=" + str(self.socket_port), ] logger.info("紫鸟客户端启动命令:%s", " ".join(cmd)) for retry_count in range(CLIENT_START_RETRIES): logger.info("第 %s/%s 次尝试启动紫鸟客户端", retry_count + 1, CLIENT_START_RETRIES) subprocess.Popen(cmd) if self._wait_until_client_ready(CLIENT_READY_TIMEOUT): logger.info("紫鸟客户端启动成功,第 %s 次尝试", retry_count + 1) time.sleep(CLIENT_POST_START_DELAY) self.update_core() return logger.warning("第 %s 次尝试启动失败,10秒内未检测到客户端启动", retry_count + 1) logger.error("紫鸟客户端启动失败,已重试 %s 次", CLIENT_START_RETRIES) raise RuntimeError(f"客户端启动失败:重试 {CLIENT_START_RETRIES} 次后仍未成功启动") def open_shop(self, shop_name: str): """打开指定店铺并完成 IP 检测. Args: shop_name: 店铺名称. Returns: 成功时返回浏览器实例.店铺不存在时返回错误信息. """ logger.info("开始打开店铺:%s", shop_name) self.start_client() self.store_id = self._find_store_oauth(self.get_browser_list(), shop_name) if not self.store_id: logger.warning("店铺不存在:%s", shop_name) return "店铺不存在" ret_json = self.open_store(self.store_id) self.store_id = ret_json.get("browserOauth") if self.store_id is None: self.store_id = ret_json.get("browserId") self.browser = self.get_browser(ret_json.get("debuggingPort")) ip_check_url = ret_json.get("ipDetectionPage") if not ip_check_url: logger.error("ip检测页地址为空,准备关闭店铺:%s", shop_name) self.close_store(self.store_id) raise RuntimeError("没有IP检测地址,为了店铺安全不打开店铺") if self.open_ip_check(self.browser, ip_check_url): logger.info("IP检测通过,打开店铺平台主页:%s", shop_name) self.open_launcher_page(ret_json.get("launcherPage"), self.browser) return self.browser logger.error("IP检测不通过,停止打开店铺:%s", shop_name) raise RuntimeError("IP检测不通过,可能是因为网络环境变化导致的,为了店铺安全不打开店铺") def close_store(self, browser_oauth=None): """关闭紫鸟店铺浏览器. Args: browser_oauth: 店铺 OAuth 标识. 未提供时使用当前打开的店铺. Returns: 紫鸟客户端返回的关闭结果. """ if browser_oauth is None: browser_oauth = self.store_id payload = self._build_payload( "stopBrowser", { "duplicate": 0, "browserOauth": browser_oauth, }, ) request_id = payload["requestId"] logger.info("开始关闭紫鸟店铺,browserOauth=%s,requestId=%s", browser_oauth, request_id) result = self._post_client(payload) status_code = self._status_code(result) if status_code == STATUS_OK: logger.info("紫鸟店铺关闭成功,browserOauth=%s", browser_oauth) return result if status_code == STATUS_LOGIN_FAILED: logger.error("关闭紫鸟店铺登录失败:%s", self._result_text(result)) raise RuntimeError(f"[close_store]登录失败 {self._result_text(result)}") logger.error("关闭紫鸟店铺失败:%s", self._result_text(result)) raise RuntimeError(f"[close_store]失败: {self._result_text(result)} ") def open_launcher_page(self, launcher_page: str, browser: Chromium = None): """打开店铺平台启动页. Args: launcher_page: 要打开的启动页 URL. browser: 浏览器实例. 未提供时使用当前浏览器实例. """ if browser is None: browser = self.browser logger.info("打开店铺平台启动页:%s", launcher_page) tab = browser.new_tab(url=launcher_page) self.tab = tab return tab def open_ip_check(self, browser: Chromium, ip_check_url: str): """打开 IP 检测页并判断网络环境是否通过. Args: browser: DrissionPage 浏览器会话. ip_check_url: IP 检测页地址. Returns: IP 检测通过时返回 True,否则返回 False. """ try: logger.info("开始打开IP检测页:%s", ip_check_url) tab = browser.latest_tab tab.get(ip_check_url) success_button = tab.ele( (By.XPATH, '//button[contains(@class, "styles_btn--success")]'), timeout=IP_CHECK_TIMEOUT, ) if success_button: logger.info("IP检测成功") return True logger.warning("IP检测超时或未找到成功按钮") return False except Exception: logger.exception("IP检测异常") return False class AmamzonBase(ZiniaoDriver): """亚马逊页面操作基类.""" def SwitchingCountries(self, country_name: str): """切换亚马逊账号国家并验证结果. 流程包括读取当前国家,打开国家下拉框,选择目标国家并等待页面加载. Args: country_name: 目标国家名称. Returns: 切换成功返回 True,否则返回 False. """ try: if self.browser is None: logger.warning("浏览器实例不存在,请先打开店铺") return False self.tab.wait.doc_loaded(timeout=COUNTRY_INITIAL_LOAD_TIMEOUT, raise_err=False) current_country = self._get_current_country(timeout=COUNTRY_LOOKUP_TIMEOUT) if current_country is None: logger.warning("无法获取当前国家信息") return False if current_country == country_name: logger.info("当前已经是目标国家 %s,无需切换", country_name) return True if not self._open_country_dropdown(): return False if not self._select_country(country_name): return False logger.info("等待国家切换后页面加载") self.tab.wait.doc_loaded() return self._verify_country(country_name) except Exception: logger.exception("切换国家时发生异常") return False def _get_current_country(self, *, timeout: int) -> Optional[str]: logger.info("正在检查当前国家") current_country_ele = self.tab.ele(COUNTRY_LABEL_XPATH, timeout=timeout) if not current_country_ele: return None current_country = current_country_ele.text.strip() logger.info("当前国家:%s", current_country) return current_country def _open_country_dropdown(self) -> bool: logger.info("正在打开国家切换下拉框") dropdown_header = self.tab.ele(COUNTRY_DROPDOWN_XPATH, timeout=COUNTRY_CLICK_TIMEOUT) if not dropdown_header: logger.warning("找不到国家切换下拉框") return False dropdown_header.click() time.sleep(COUNTRY_DROPDOWN_DELAY) logger.info("正在展开国家列表") first_item = self.tab.ele(COUNTRY_LIST_ITEM_XPATH, timeout=COUNTRY_CLICK_TIMEOUT) if not first_item: logger.warning("找不到国家列表项") return False first_item.click() time.sleep(COUNTRY_DROPDOWN_DELAY) return True def _select_country(self, country_name: str) -> bool: logger.info("正在切换到国家:%s", country_name) target_country = self.tab.ele(self._country_option_xpath(country_name), timeout=COUNTRY_CLICK_TIMEOUT) if not target_country: logger.warning("找不到目标国家:%s", country_name) return False target_country.click() return True @staticmethod def _country_option_xpath(country_name: str) -> str: return ( "xpath://div[@class=" '"dropdown-account-switcher-list-item dropdown-account-switcher-list-item-indented" ' f'and @title="{country_name}"]' ) def _verify_country(self, country_name: str) -> bool: new_country = self._get_current_country(timeout=COUNTRY_CLICK_TIMEOUT) if new_country is None: logger.warning("无法验证国家切换结果") return False if new_country == country_name: logger.info("国家切换成功:%s", new_country) return True logger.warning("国家切换失败,当前国家:%s,目标国家:%s", new_country, country_name) return False def need_login(self): """判断当前页面是否需要登录.""" time.sleep(NEED_LOGIN_DELAY) self.tab.wait.doc_loaded(timeout=DOC_LOAD_TIMEOUT, raise_err=False) need_login_ele = self.tab.eles(NEED_LOGIN_XPATH, timeout=NEED_LOGIN_TIMEOUT) if len(need_login_ele) > 0: logger.info("检测到需要登录元素") return True logger.info("未检测到登录元素") return False def login(self, password, username=""): """登录当前页面,必要时等待并提交一次性验证码. Args: password: 登录密码. username: 兼容保留参数,当前未使用. Returns: 登录成功返回 True,失败返回 False. """ try: self.tab.wait.doc_loaded(timeout=DOC_LOAD_TIMEOUT, raise_err=False) for _ in range(LOGIN_ATTEMPTS): self._submit_password_if_present(password) self.tab.wait.doc_loaded(timeout=DOC_LOAD_TIMEOUT, raise_err=False) if self._send_otp_if_present(): continue if self._submit_ready_otp(): return True self._click_final_login_submit_if_present() except Exception: logger.exception("登录过程中发生异常") return False def _submit_password_if_present(self, password) -> None: pwd_input = self.tab.eles(PASSWORD_INPUT_XPATH, timeout=PASSWORD_INPUT_TIMEOUT) if len(pwd_input) > 0: logger.info("检测到密码输入框,准备输入密码") pwd_input[0].input(password, clear=True) submit_btn = self.tab.eles(PASSWORD_SUBMIT_XPATH, timeout=PASSWORD_SUBMIT_TIMEOUT) if len(submit_btn) > 0: logger.info("点击登录提交按钮") submit_btn[0].click() def _send_otp_if_present(self) -> bool: send_code = self.tab.eles(OTP_SEND_XPATH, timeout=OTP_SEND_TIMEOUT) logger.info("发送一次性密码元素数量:%s", len(send_code)) if len(send_code) == 0: return False logger.info("检测到发送一次性密码按钮,准备点击") send_code[0].click() time.sleep(OTP_SEND_DELAY) self.tab.wait.doc_loaded(timeout=DOC_LOAD_TIMEOUT, raise_err=False) return True def _submit_ready_otp(self) -> bool: for _ in range(OTP_INPUT_LOOKUP_ATTEMPTS): otp_input = self._find_otp_input() if otp_input is None: return False if self._wait_for_otp_value(otp_input) and self._submit_otp_if_ready(): return True return False def _find_otp_input(self): otp_code_input = self.tab.eles(OTP_INPUT_XPATH, timeout=OTP_INPUT_TIMEOUT) logger.info("验证码输入框数量:%s", len(otp_code_input)) if len(otp_code_input) == 0: return None otp_code_input[0].wait.displayed(timeout=OTP_INPUT_DISPLAY_TIMEOUT, raise_err=False) return otp_code_input[0] @staticmethod def _has_otp_value(otp_input) -> bool: return otp_input.value is not None and otp_input.value.strip() != "" def _wait_for_otp_value(self, otp_input) -> bool: for _ in range(OTP_CODE_WAIT_SECONDS): if self._has_otp_value(otp_input): logger.info("检测到验证码输入完成") return True time.sleep(1) return False def _submit_otp_if_ready(self) -> bool: submit_btn = self.tab.eles(OTP_SUBMIT_XPATH, timeout=OTP_SUBMIT_TIMEOUT) if len(submit_btn) == 0: return False submit_btn[0].click() self.tab.wait.doc_loaded(timeout=OTP_RESULT_LOAD_TIMEOUT, raise_err=False) error_mes = self.tab.eles(OTP_ERROR_XPATH, timeout=OTP_ERROR_TIMEOUT) if len(error_mes) > 0: logger.warning("验证码输入错误:%s", error_mes[0].text) self.tab.refresh() return False logger.info("登录成功") return True def _click_final_login_submit_if_present(self) -> None: submit_btn = self.tab.eles(OTP_SUBMIT_XPATH, timeout=LOGIN_FINAL_SUBMIT_TIMEOUT) if len(submit_btn) > 0: logger.info("点击二次登录提交按钮") submit_btn[0].click()