import json import os import subprocess import time import uuid from typing import Any, Dict, List, Literal, Optional, TypedDict from loguru import logger 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 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={}", 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("正在从注册表读取紫鸟客户端路径:{}\\{}", 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("已获取紫鸟客户端路径:{}", exe_path) return exe_path except FileNotFoundError: logger.warning("注册表中未找到紫鸟客户端路径:{}\\{}", root_name, key_path) logger.error("未能获取紫鸟客户端路径,protocol_name={}", 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("更新内核返回:{}", result) if self._handle_update_core_result(result): return if result is None: continue logger.info("等待更新内核完成:{}", 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") self._kill_process_by_name("starter.exe") else: process_name = "ziniao.exe" logger.info("准备结束紫鸟客户端进程:{}", process_name) self._kill_process_by_name(process_name) time.sleep(PROCESS_KILL_DELAY) def _kill_process_by_name(self, process_name: str) -> None: result = subprocess.run( ["taskkill", "/f", "/t", "/im", process_name], capture_output=True, text=True, encoding="utf-8", errors="ignore", check=False, ) output = (result.stdout or result.stderr or "").strip() if result.returncode == 0: if output: logger.info("结束进程返回:{}", output) return if "没有找到进程" in output or "not found" in output.lower(): logger.info("进程未运行,跳过结束:{}", process_name) return logger.warning("结束进程失败:{} | {}", process_name, output or f"returncode={result.returncode}") def get_browser_list(self) -> Optional[list[Dict[str, Any]]]: """获取紫鸟浏览器店铺列表. Returns: 浏览器店铺列表.登录失败或请求失败时返回 None. """ payload = self._build_payload("getBrowserList") request_id = payload["requestId"] logger.info("开始获取紫鸟浏览器列表,requestId={}", 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("获取紫鸟浏览器列表成功,数量={}", len(browser_list or [])) return browser_list if status_code == STATUS_LOGIN_FAILED: logger.error("获取紫鸟浏览器列表登录失败:{}", self._result_text(result)) return None logger.error("获取紫鸟浏览器列表失败:{}", 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={},requestId={}", store_info, request_id) result = self._post_client(payload) status_code = self._status_code(result) if status_code == STATUS_OK: logger.info("紫鸟店铺打开成功,debuggingPort={}", result.get("debuggingPort")) return result if status_code == STATUS_LOGIN_FAILED: logger.error("打开紫鸟店铺登录失败:{}", self._result_text(result)) raise RuntimeError(f"[open_store]登录失败 {self._result_text(result)}") logger.error("打开紫鸟店铺失败:{}", 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={}", port) browser = Chromium(port) logger.info("DrissionPage 浏览器连接完成,port={}", port) return browser def start_client(self): """启动紫鸟客户端并更新浏览器内核.""" if self._is_client_ready(): logger.info("端口 {} 已启动,跳过启动紫鸟客户端", self.socket_port) return logger.info("端口 {} 未启动,开始启动紫鸟客户端", 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("紫鸟客户端启动命令:{}", " ".join(cmd)) for retry_count in range(CLIENT_START_RETRIES): logger.info("第 {}/{} 次尝试启动紫鸟客户端", retry_count + 1, CLIENT_START_RETRIES) subprocess.Popen(cmd) if self._wait_until_client_ready(CLIENT_READY_TIMEOUT): logger.info("紫鸟客户端启动成功,第 {} 次尝试", retry_count + 1) time.sleep(CLIENT_POST_START_DELAY) self.update_core() return logger.warning("第 {} 次尝试启动失败,10秒内未检测到客户端启动", retry_count + 1) logger.error("紫鸟客户端启动失败,已重试 {} 次", CLIENT_START_RETRIES) raise RuntimeError(f"客户端启动失败:重试 {CLIENT_START_RETRIES} 次后仍未成功启动") def open_shop(self, shop_name: str): """打开指定店铺并完成 IP 检测. Args: shop_name: 店铺名称. Returns: 成功时返回浏览器实例.店铺不存在时返回错误信息. """ logger.info("开始打开店铺:{}", 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("店铺不存在:{}", 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检测页地址为空,准备关闭店铺:{}", shop_name) self.close_store(self.store_id) raise RuntimeError("没有IP检测地址,为了店铺安全不打开店铺") if self.open_ip_check(self.browser, ip_check_url): logger.info("IP检测通过,打开店铺平台主页:{}", shop_name) self.open_launcher_page(ret_json.get("launcherPage"), self.browser) return self.browser logger.error("IP检测不通过,停止打开店铺:{}", 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={},requestId={}", browser_oauth, request_id) result = self._post_client(payload) status_code = self._status_code(result) if status_code == STATUS_OK: logger.info("紫鸟店铺关闭成功,browserOauth={}", browser_oauth) return result if status_code == STATUS_LOGIN_FAILED: logger.error("关闭紫鸟店铺登录失败:{}", self._result_text(result)) raise RuntimeError(f"【close_store】登录失败 {self._result_text(result)}") logger.error("关闭紫鸟店铺失败:{}", 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("打开店铺平台启动页:{}", 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检测页:{}", 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 switch_to_country(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("当前已经是目标国家 {},无需切换", 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("当前国家:{}", 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("正在切换到国家:{}", country_name) target_country = self.tab.ele(self._country_option_xpath(country_name), timeout=COUNTRY_CLICK_TIMEOUT) if not target_country: logger.warning("找不到目标国家:{}", 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("国家切换成功:{}", new_country) return True logger.warning("国家切换失败,当前国家:{},目标国家:{}", 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("发送一次性密码元素数量:{}", 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("验证码输入框数量:{}", 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("验证码输入错误:{}", 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()