From e478ac4d094ca716b7cf97eeeb530f337962f8f2 Mon Sep 17 00:00:00 2001 From: yml2213 Date: Wed, 24 Jun 2026 15:37:48 +0800 Subject: [PATCH] =?UTF-8?q?=E7=AE=80=E5=8C=96=E7=99=BB=E5=BD=95=E6=B5=81?= =?UTF-8?q?=E7=A8=8B=EF=BC=9A=E7=A0=8D=E6=8E=89=E4=BB=A3=E7=90=86=E6=B1=A0?= =?UTF-8?q?=EF=BC=8C=E6=9E=81=E9=AA=8C=E4=B8=8D=E5=86=85=E9=83=A8=E5=88=87?= =?UTF-8?q?=E4=BB=A3=E7=90=86?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 新建 ProxyFetcher 替代 ProxyManager,每次从API取1个代理,无池无锁无冷却 - 极验最多3次尝试,失败直接抛异常回到login整体重试取新代理 - 删掉极验内部的 refresh_proxy/proxy_switches/_soft_fail_streak 逻辑 - login整体重试时取新代理从头走完整流程 - 代理验证超时从(5,8)缩短为(3,5),白名单同步后立即重试不等待 - login_service 删掉60行代理池预检/等待代码 - 修复邮件验证码日期比较:Roundcube分钟精度与秒级时间戳对齐 --- core/douyu/login.py | 247 +++++++------------------- core/douyu/proxy.py | 3 + core/douyu/proxy_fetcher.py | 69 +++++++ core/douyu/proxy_manager.py | 13 +- core/douyu/proxy_resolver.py | 4 +- core/douyu/proxy_verifier.py | 7 +- web/backend/services/login_service.py | 97 ++-------- 7 files changed, 162 insertions(+), 278 deletions(-) create mode 100644 core/douyu/proxy_fetcher.py diff --git a/core/douyu/login.py b/core/douyu/login.py index 54950ec..85b4d82 100644 --- a/core/douyu/login.py +++ b/core/douyu/login.py @@ -12,7 +12,7 @@ from loguru import logger from .cookie_enricher import CookieEnricher from .crypto import encrypt_password, encrypt_nickname_or_phone from .email_verifier import EmailVerifier -from .proxy import ProxyManager, get_proxy_manager +from .proxy_fetcher import ProxyFetcher from core.geetest.v3_slide.solver import ( _generate_seed, get_w1, get_w2, @@ -80,31 +80,34 @@ class DouyuLogin: whitelist_ukey: str = "", whitelist_platform: str = "xiequ", whitelist_credentials: dict = None, - proxy_manager: Optional[ProxyManager] = None, + proxy_fetcher: Optional[ProxyFetcher] = None, stop_event: Optional[threading.Event] = None, + # 向后兼容:旧代码传 proxy_manager 时自动转为 proxy_fetcher + proxy_manager=None, ): self.account = account self.proxy = proxy self.timeout = timeout self.max_proxy_retries = max_proxy_retries # 0=无限切换直到成功 - self.max_login_retries = max_login_retries # 0=无限整体重试直到成功 + self.max_login_retries = max_login_retries # 0=无限重试直到成功 self.max_total_time = max_total_time # 0=不限制单账号登录总时长 self.stop_event = stop_event self.session = requests.Session() - # 初始化代理管理器(优先使用外部传入的共享实例,避免并发刷新冲突) - if proxy_manager: - self.proxy_manager = proxy_manager + # 初始化代理获取器 + if proxy_fetcher: + self.proxy_fetcher = proxy_fetcher elif proxy_api_url: - self.proxy_manager = get_proxy_manager( - proxy_api_url, - whitelist_uid=whitelist_uid, - whitelist_ukey=whitelist_ukey, + self.proxy_fetcher = ProxyFetcher( + api_url=proxy_api_url, whitelist_platform=whitelist_platform, - whitelist_credentials=whitelist_credentials, + whitelist_credentials=whitelist_credentials or ( + {"uid": whitelist_uid, "ukey": whitelist_ukey} + if whitelist_uid and whitelist_ukey else None + ), ) else: - self.proxy_manager = None + self.proxy_fetcher = None self._current_proxy_url: Optional[str] = None self._cookie_enrich_error = "" @@ -150,7 +153,7 @@ class DouyuLogin: self._apply_proxy() def _apply_proxy(self, proxy: str = None) -> None: - """应用代理到Session,并记录当前代理URL供刷新时mark_bad""" + """应用代理到Session""" if proxy: # 使用指定的代理 self.session.proxies = { @@ -159,7 +162,7 @@ class DouyuLogin: } self._current_proxy_url = proxy elif self.proxy: - # 使用配置的代理 + # 使用配置的静态代理 if isinstance(self.proxy, str): self.session.proxies = { 'http': self.proxy, @@ -173,75 +176,17 @@ class DouyuLogin: if url } self._current_proxy_url = self.proxy.get('http') or self.proxy.get('https') - else: - # 从代理管理器获取代理 - if self.proxy_manager: - new_proxy = self.proxy_manager.get_proxy() - if new_proxy: - self.session.proxies = { - 'http': new_proxy, - 'https': new_proxy, - } - self._current_proxy_url = new_proxy - - def _refresh_proxy(self, mark_bad: bool = True) -> Optional[str]: - """ - 刷新代理IP。 - - Args: - mark_bad: 是否标记当前代理为坏。默认True(代理确实不可用时)。 - 设为False时仅换代理,不标记坏(临时网络波动,代理本身可能没问题)。 - """ - if not self.proxy_manager: - return None - - if mark_bad and self._current_proxy_url: - self.proxy_manager.mark_bad(self._current_proxy_url) - - new_proxy = self.proxy_manager.get_proxy() - self._ensure_not_stopped() - if new_proxy: - self._apply_proxy(new_proxy) - logger.info(f"已切换代理: {new_proxy}{' (旧代理已标记坏)' if mark_bad else ' (旧代理保留)'}") - return new_proxy - - @staticmethod - def _is_proxy_connection_error(err_str: str) -> bool: - """判断异常是否为代理连接类错误(代理已死),而非临时网络波动。""" - err_lower = err_str.lower() - # 代理连接/超时类:代理本身不可达 - proxy_dead_keywords = [ - 'proxyerror', - 'tunnel connection failed', - 'connecttimeouterror', - 'connection refused', - 'unable to connect to proxy', - 'proxy connection', - '502 bad gateway', - '503 service unavailable', - ] - if any(kw in err_lower for kw in proxy_dead_keywords): - return True - # SSL/证书错误也视为代理问题 - if 'ssl' in err_lower and 'proxy' in err_lower: - return True - return False - - @staticmethod - def _is_geetest_network_error(err_str: str) -> bool: - """ - 判断异常是否为极验接口的网络类错误(代理到极验不可达/被限流)。 - 这类错误说明当前代理已被极验识别或限流,应立即换代理而非软重试。 - """ - err_lower = err_str.lower() - network_error_keywords = [ - '网络不给力', # 极验 get.php 返回的网络错误 - 'connection aborted', # 远程断开连接 - 'remotedisconnected', # requests RemoteDisconnected - 'read timed out', # 读取超时(极验接口) - 'api.geetest.com', # 极验接口相关错误 - ] - return any(kw in err_lower for kw in network_error_keywords) + elif self.proxy_fetcher: + # 从代理获取器取新代理 + new_proxy = self.proxy_fetcher.fetch_new_proxy() + if new_proxy: + self.session.proxies = { + 'http': new_proxy, + 'https': new_proxy, + } + self._current_proxy_url = new_proxy + else: + logger.warning("获取代理失败,将尝试直连") @staticmethod def _truncate_error(err_str: str, max_len: int = 80) -> str: @@ -255,10 +200,10 @@ class DouyuLogin: parsed = urlsplit(url) return urlunsplit((parsed.scheme, parsed.netloc, parsed.path, "", "")) - def _request(self, method: str, url: str, max_retries: int = 3, **kwargs) -> requests.Response: + def _request(self, method: str, url: str, max_retries: int = 2, **kwargs) -> requests.Response: """ 统一发送请求,附带分段超时和更明确的错误信息。 - 代理连接失败时自动重试获取新的代理IP。 + 代理连接失败时直接抛异常,由 login 整体重试换新代理。 """ timeout = kwargs.pop('timeout', self.timeout) safe_url = self._safe_url(url) @@ -285,22 +230,12 @@ class DouyuLogin: elapsed = time.monotonic() - started err_str = str(exc) is_proxy_err = "proxy" in err_str.lower() or "Proxy" in type(exc).__name__ - if is_proxy_err: - logger.warning(f"代理连接失败,尝试 {attempt + 1}/{max_retries}: {exc}") - if attempt < max_retries - 1: - self._refresh_proxy() - self._sleep_interruptible(1) - continue - raise ConnectionError( - f"{method.upper()} {safe_url} 代理连接失败,已重试 {max_retries} 次" - ) from exc - # 非代理的 ConnectionError 也重试 - logger.warning(f"连接失败,尝试 {attempt + 1}/{max_retries}: {exc}") - if attempt < max_retries - 1: + if is_proxy_err and attempt < max_retries - 1: + logger.warning(f"代理连接失败,尝试 {attempt + 1}/{max_retries}: {self._truncate_error(err_str)}") self._sleep_interruptible(1) continue raise ConnectionError( - f"{method.upper()} {safe_url} 连接失败,已重试 {max_retries} 次" + f"{method.upper()} {safe_url} 连接失败: {self._truncate_error(err_str)}" ) from exc except requests.RequestException as exc: elapsed = time.monotonic() - started @@ -308,6 +243,8 @@ class DouyuLogin: f"{method.upper()} {safe_url} 请求失败,耗时 {elapsed:.1f}s: {exc}" ) from exc + raise ConnectionError(f"{method.upper()} {safe_url} 连接失败,已重试 {max_retries} 次") + def _request_json(self, method: str, url: str, source: str, **kwargs) -> dict: """请求 JSON 接口,并在响应异常时输出可定位的信息。""" response = self._request(method, url, **kwargs) @@ -327,7 +264,7 @@ class DouyuLogin: """ 完整登录流程(带整体重试)。 - 任何步骤失败时,换新代理从头重跑。 + 每一轮用一个代理走完所有步骤,任何步骤失败就取新代理从头重来。 max_login_retries=0 表示无限重试,max_total_time=0 表示不限制总时长。 Returns: @@ -350,10 +287,10 @@ class DouyuLogin: if attempt > 1: if self.max_login_retries > 0: - logger.info(f"登录整体重试 {attempt}/{self.max_login_retries},换代理重新开始") + logger.info(f"登录整体重试 {attempt}/{self.max_login_retries},换新代理从头开始") else: - logger.info(f"登录整体重试 {attempt} (无限重试),换代理重新开始") - # 重试前:换新代理 + 重置 session(清 cookies) + logger.info(f"登录整体重试 {attempt} (无限重试),换新代理从头开始") + # 重试前:取新代理 + 重置 session self._prepare_retry() try: @@ -390,9 +327,6 @@ class DouyuLogin: message = f"登录成功,补CK失败: {self._cookie_enrich_error}" logger.success(f"登录成功! Cookie长度: {len(cookie)}") - # 成功后归还代理到池,让其他账号复用 - if self.proxy_manager and self._current_proxy_url: - self.proxy_manager.release_proxy(self._current_proxy_url) return LoginResult(success=True, cookie=cookie, message=message) except InterruptedError as e: @@ -408,24 +342,19 @@ class DouyuLogin: has_retry = self.max_login_retries <= 0 or attempt < self.max_login_retries has_time = self.max_total_time <= 0 or elapsed < self.max_total_time if has_retry and has_time: - # 还有重试机会且未超时:标记当前代理坏,下一轮自动换新代理 - if self.proxy_manager and self._current_proxy_url: - self.proxy_manager.mark_bad(self._current_proxy_url) - self._sleep_interruptible(2) + self._sleep_interruptible(1) continue # 所有重试耗尽或超时 return LoginResult(success=False, message=str(e)) def _prepare_retry(self) -> None: - """重试前准备:换新代理、重置 session cookies。""" - # 重置 session(清掉旧 cookies,避免残留状态干扰) + """重试前准备:取新代理、重置 session cookies。""" self.session.cookies.clear() - # 刷新代理 - if self.proxy_manager: - new_proxy = self.proxy_manager.get_proxy() + if self.proxy_fetcher: + new_proxy = self.proxy_fetcher.fetch_new_proxy() if new_proxy: self._apply_proxy(new_proxy) - # 如果获取新代理失败,保留旧代理继续尝试 + logger.info(f"重试换新代理: {new_proxy}") def _first_login(self) -> Tuple[str, str, str, dict]: """ @@ -482,12 +411,9 @@ class DouyuLogin: def _solve_geetest(self, gt: str, challenge: str, deadline: float = 0) -> Tuple[str, str]: """ - 解决极验 fullpage 验证(带重试机制) + 解决极验 fullpage 验证(最多3次尝试,失败直接抛异常回到login换新代理) - 区分三类错误: - - 代理死亡(ProxyError、连接超时):立即换代理 - - 网络类错误(网络不给力、Connection aborted):立即换代理 - - 极验逻辑失败(slide等):原代理重试,连续2次才换 + 不在内部做代理切换——短效代理寿命宝贵,换代理由 login 整体重试负责。 Args: gt: 极验gt参数 @@ -509,48 +435,27 @@ class DouyuLogin: _geetest_semaphore.release() def _solve_geetest_inner(self, gt: str, challenge: str, deadline: float = 0) -> Tuple[str, str]: - """极验验证内部实现(已获取并发信号量)。""" - max_proxy_switches = self.max_proxy_retries if self.max_proxy_retries > 0 else 999999 - proxy_switches = 0 - - # ── 极验验证最大尝试次数(超出后回到 login 整体重试)── - _MAX_GEETEST_ATTEMPTS = 5 - - # 连续临时失败计数(同一代理下),超过阈值才换代理 - _soft_fail_streak = 0 - _SOFT_FAIL_THRESHOLD = 2 - - def refresh_proxy(mark_bad: bool) -> None: - """按代理切换上限刷新代理。""" - nonlocal proxy_switches - if not self.proxy_manager: - return - if proxy_switches >= max_proxy_switches: - raise ValueError(f"极验验证代理切换次数已达上限 {self.max_proxy_retries}") - new_proxy = self._refresh_proxy(mark_bad=mark_bad) - if new_proxy: - proxy_switches += 1 + """极验验证内部实现:最多3次尝试,失败直接抛异常。""" + _MAX_ATTEMPTS = 3 attempt = 0 while True: attempt += 1 self._ensure_not_stopped() - # ── 极验验证尝试上限 ── - if attempt > _MAX_GEETEST_ATTEMPTS: + if attempt > _MAX_ATTEMPTS: raise ValueError( - f"极验验证已尝试 {attempt} 次,超过上限 {_MAX_GEETEST_ATTEMPTS}," + f"极验验证已尝试 {attempt - 1} 次,超过上限 {_MAX_ATTEMPTS}," f"回到 login 整体重试换新代理" ) - # 超时兜底:极验验证不应超过登录整体时间上限 + # 超时兜底 if deadline and time.monotonic() > deadline: raise ValueError(f"极验验证超时(登录整体时间耗尽)") try: - logger.info(f"极验验证尝试 {attempt}/{_MAX_GEETEST_ATTEMPTS}") + logger.info(f"极验验证尝试 {attempt}/{_MAX_ATTEMPTS}") - # 按斗鱼登录页 HAR:fullpage 智能检测流程,不进入图片滑块。 str_16 = _generate_seed() proxies = dict(self.session.proxies) @@ -581,16 +486,8 @@ class DouyuLogin: logger.success(f"极验 fullpage 验证成功! validate={validate[:20]}...") return validate, seccode else: - # 极验返回失败(slide 等)—— 通常是临时问题,先原代理重试 - _soft_fail_streak += 1 - if _soft_fail_streak >= _SOFT_FAIL_THRESHOLD: - logger.warning(f"极验验证失败: {message},同一代理连续 {_soft_fail_streak} 次,换代理") - refresh_proxy(mark_bad=False) - _soft_fail_streak = 0 - else: - logger.warning(f"极验验证失败: {message},原代理重试 ({_soft_fail_streak}/{_SOFT_FAIL_THRESHOLD})") - self._sleep_interruptible(1) - continue + logger.warning(f"极验验证失败: {message},回到 login 换新代理") + raise ValueError(f"极验验证失败: {message}") else: validate = str(result) if validate: @@ -598,36 +495,14 @@ class DouyuLogin: logger.success(f"极验 fullpage 验证成功! validate={validate[:20]}...") return validate, seccode + except ValueError: + # 极验逻辑失败(slide等)或超过上限,直接上抛 + raise except Exception as e: err_str = str(e) - if "代理切换次数已达上限" in err_str or "超过上限" in err_str: - raise - is_proxy_dead = self._is_proxy_connection_error(err_str) - is_network_error = self._is_geetest_network_error(err_str) - - if is_proxy_dead: - # 代理确实不可用:立即换,标记坏 - logger.warning(f"极验验证代理连接失败: {self._truncate_error(err_str)},换代理") - refresh_proxy(mark_bad=True) - _soft_fail_streak = 0 - elif is_network_error: - # 网络类错误(极验限流/网络不给力/Connection aborted):立即换代理 - logger.warning(f"极验网络错误: {self._truncate_error(err_str)},立即换代理") - refresh_proxy(mark_bad=False) - _soft_fail_streak = 0 - else: - # 其他临时异常(KeyError 等):先原代理重试 - _soft_fail_streak += 1 - if _soft_fail_streak >= _SOFT_FAIL_THRESHOLD: - logger.warning(f"极验验证临时异常: {self._truncate_error(err_str)},连续 {_soft_fail_streak} 次,换代理") - refresh_proxy(mark_bad=False) - _soft_fail_streak = 0 - else: - logger.warning(f"极验验证临时异常: {self._truncate_error(err_str)},原代理重试 ({_soft_fail_streak}/{_SOFT_FAIL_THRESHOLD})") - self._sleep_interruptible(2) - continue - - raise ValueError("极验验证失败(无限重试模式仍未能通过)") + # 所有网络/代理异常都直接上抛,由 login 换新代理 + logger.warning(f"极验验证异常: {self._truncate_error(err_str)},回到 login 换新代理") + raise ValueError(f"极验验证异常: {self._truncate_error(err_str)}") from e def _second_login(self, gt: str, challenge: str, validate: str, seccode: str, code_token: str) -> str: @@ -647,8 +522,6 @@ class DouyuLogin: encrypted_username = encrypt_nickname_or_phone(self.account.username) encrypted_password = encrypt_password(self.account.password) - # 参考HAR文件中的完整参数 - # 注意:geetest_challenge应该使用第一次登录返回的challenge data = { 'type': '1', 'nicknameOrPhoneEncrypt': encrypted_username, diff --git a/core/douyu/proxy.py b/core/douyu/proxy.py index bb65408..c3521f9 100644 --- a/core/douyu/proxy.py +++ b/core/douyu/proxy.py @@ -1,13 +1,16 @@ """代理管理模块兼容导出。""" +from .proxy_fetcher import ProxyFetcher, get_proxy_fetcher from .proxy_manager import ProxyManager, get_proxy_manager from .proxy_parser import parse_proxy_response from .proxy_resolver import ProxyResolver, resolve_working_proxy from .proxy_verifier import verify_proxies_concurrent, verify_proxy_url __all__ = [ + "ProxyFetcher", "ProxyManager", "ProxyResolver", + "get_proxy_fetcher", "get_proxy_manager", "parse_proxy_response", "resolve_working_proxy", diff --git a/core/douyu/proxy_fetcher.py b/core/douyu/proxy_fetcher.py new file mode 100644 index 0000000..578dc1a --- /dev/null +++ b/core/douyu/proxy_fetcher.py @@ -0,0 +1,69 @@ +"""简化版代理获取器(替代 ProxyManager)。 + +专为短效代理设计:无池、无锁、无冷却、无复用。 +每次调用 fetch_new_proxy() 从 API 取 1 个新代理,用完即弃。 + +线程安全:白名单同步由 ProxyResolver 内部的锁保证。 +""" + +from typing import Optional + +from loguru import logger + +from .proxy_resolver import ProxyResolver +from .proxy_whitelist import DouyuWhitelistSyncer + + +class ProxyFetcher: + """每次从代理 API 取 1 个新代理,无缓存无复用。""" + + def __init__( + self, + api_url: str, + whitelist_platform: str = "xiequ", + whitelist_credentials: dict = None, + ): + self.api_url = api_url + + _wl_platform = whitelist_platform or "xiequ" + _wl_credentials = whitelist_credentials + + self._whitelist_syncer = ( + DouyuWhitelistSyncer(platform=_wl_platform, credentials=_wl_credentials) + if _wl_credentials + else None + ) + + def fetch_new_proxy(self, max_attempts: int = 3) -> Optional[str]: + """从代理 API 获取 1 个可用代理,失败返回 None。""" + resolver = ProxyResolver( + api_url=self.api_url, + whitelist_syncer=self._whitelist_syncer, + sync_local_exit_ip=bool(self._whitelist_syncer), + sync_whitelist_once=False, + ) + result, msg = resolver.fetch_verified( + max_attempts=max_attempts, + return_all=False, + ) + if isinstance(result, str): + logger.info(f"获取新代理: {result}") + return result + if isinstance(result, list) and result: + logger.info(f"获取新代理: {result[0]}") + return result[0] + logger.warning(f"获取代理失败: {msg}") + return None + + +def get_proxy_fetcher( + api_url: str, + whitelist_platform: str = "xiequ", + whitelist_credentials: dict = None, +) -> ProxyFetcher: + """获取代理获取器实例。""" + return ProxyFetcher( + api_url, + whitelist_platform=whitelist_platform, + whitelist_credentials=whitelist_credentials, + ) diff --git a/core/douyu/proxy_manager.py b/core/douyu/proxy_manager.py index 9cc464b..99566a4 100644 --- a/core/douyu/proxy_manager.py +++ b/core/douyu/proxy_manager.py @@ -155,18 +155,23 @@ class ProxyManager: return None def _fetch_and_verify_all(self, max_attempts: int) -> tuple[Optional[list[str]], str]: - """调代理 API 获取一批代理,并发验证所有可用代理。""" + """调代理 API 获取一批代理,并发验证。 + + 对短效代理友好:拿到第一个可用代理就立即返回, + 不再等所有代理验证完毕(验证期间短效代理会过期浪费)。 + """ resolver = ProxyResolver( api_url=self.api_url, whitelist_syncer=self._whitelist_syncer, sync_local_exit_ip=bool(self._whitelist_syncer), sync_whitelist_once=False, ) - available, msg = resolver.fetch_verified(max_attempts=max_attempts, return_all=True) - if isinstance(available, list): - return available, msg + # return_all=False: 找到第一个可用就返回,节省短效代理时间 + available, msg = resolver.fetch_verified(max_attempts=max_attempts, return_all=False) if isinstance(available, str): return [available], msg + if isinstance(available, list): + return available, msg return None, msg def mark_bad(self, proxy_url: str) -> None: diff --git a/core/douyu/proxy_resolver.py b/core/douyu/proxy_resolver.py index 9d292de..ee1d626 100644 --- a/core/douyu/proxy_resolver.py +++ b/core/douyu/proxy_resolver.py @@ -72,7 +72,6 @@ class ProxyResolver: ok, sync_msg = self._sync_ip(local_ip) if ok: self._log('info', f"白名单同步成功: {sync_msg}") - time.sleep(2) else: self._log('warning', f"白名单同步失败: {sync_msg}") elif local_ip == self._last_synced_ip: @@ -130,8 +129,7 @@ class ProxyResolver: ok, sync_msg = self._sync_ip(whitelist_ip) self._log('success' if ok else 'error', f'白名单同步: {sync_msg}') if ok: - self._log('info', '白名单已更新,等待2秒后重试...') - time.sleep(2) + self._log('info', '白名单已更新,立即重试...') continue return None, f'白名单同步失败: {sync_msg}' diff --git a/core/douyu/proxy_verifier.py b/core/douyu/proxy_verifier.py index 0b4183a..28180c9 100644 --- a/core/douyu/proxy_verifier.py +++ b/core/douyu/proxy_verifier.py @@ -7,14 +7,11 @@ import requests from loguru import logger -def verify_proxy_url(proxy_url: str, timeout: tuple = (5, 8)) -> tuple[bool, str]: +def verify_proxy_url(proxy_url: str, timeout: tuple = (3, 5)) -> tuple[bool, str]: """ 验证代理是否可用,只验证斗鱼主站可达。 - 不再验证极验接口,原因: - 1. 极验接口验证耗时 5-8秒,对短效代理是巨大浪费 - 2. 极验限流导致大量代理被误判为不可用 - 3. 登录流程中极验验证本身会检测代理到极验的连通性 + 超时缩短为 (3, 5) 以减少短效代理在验证期间过期浪费。 Returns: (是否可用, 消息) diff --git a/web/backend/services/login_service.py b/web/backend/services/login_service.py index bd6c069..e57dac6 100644 --- a/web/backend/services/login_service.py +++ b/web/backend/services/login_service.py @@ -12,16 +12,13 @@ from typing import Optional from sqlalchemy.orm import Session from core.douyu import DouyuLogin -from core.douyu.proxy import resolve_working_proxy, get_proxy_manager +from core.douyu.proxy_fetcher import ProxyFetcher from ..models import Account as AccountModel, LoginTask, ProxyConfig as ProxyConfigModel class LoginBatchRunner: """批量登录执行器,在线程中运行,通过 ThreadPoolExecutor 并发登录多个账号。""" - # 代理池耗尽时等待恢复的最大秒数 - _PROXY_WAIT_MAX = 60 - def __init__( self, db: Session, @@ -52,21 +49,21 @@ class LoginBatchRunner: self._counter_lock = threading.Lock() self._completed = 0 - # 共享代理管理器(带锁,避免并发白名单限流;极验失败时可刷新代理) - self._shared_proxy_manager = None + # 共享代理获取器(无池,每次取新代理) + self._shared_proxy_fetcher = None if proxy_config and proxy_config.enabled and proxy_config.api_url: wl_platform = "xiequ" wl_credentials = None if proxy_config.whitelist_enabled: wl_platform = getattr(proxy_config, 'whitelist_platform', None) or "xiequ" wl_credentials = getattr(proxy_config, 'whitelist_credentials', None) - # 向后兼容:旧字段有值但新字段为空 + # 向后兼容 if not wl_credentials and proxy_config.whitelist_uid and proxy_config.whitelist_ukey: wl_platform = "xiequ" wl_credentials = {"uid": proxy_config.whitelist_uid, "ukey": proxy_config.whitelist_ukey} - self._shared_proxy_manager = get_proxy_manager( - proxy_config.api_url, + self._shared_proxy_fetcher = ProxyFetcher( + api_url=proxy_config.api_url, whitelist_platform=wl_platform, whitelist_credentials=wl_credentials, ) @@ -91,13 +88,7 @@ class LoginBatchRunner: ) def _resolve_static_proxy(self) -> tuple[Optional[dict], str]: - """ - 解析静态代理配置(仅处理无代理和静态代理场景)。 - - API代理由 DouyuLogin 通过 proxy_manager 内部管理, - 不在此处预先获取——登录过程中的代理切换(极验失败、整体重试) - 都在 DouyuLogin 内部自治完成。 - """ + """解析静态代理配置。""" if not self.proxy_config or not self.proxy_config.enabled: return None, '' @@ -106,11 +97,11 @@ class LoginBatchRunner: proxy_url = self.proxy_config.http or self.proxy_config.https return {'http': proxy_url, 'https': proxy_url}, f'使用静态代理: {proxy_url}' - # API代理:不在此处获取,由 DouyuLogin 通过 proxy_manager 内部管理 + # API代理:由 DouyuLogin 通过 proxy_fetcher 内部管理 return None, '' def _execute_one(self, task_id: int, acc_info: dict, total: int): - """在独立线程中执行单个账号登录,使用独立的 DB 会话和代理。""" + """在独立线程中执行单个账号登录,使用独立的 DB 会话。""" if self._stop.is_set(): self._push_log("warning", f"任务已停止,跳过: {acc_info['username']}") return @@ -135,56 +126,8 @@ class LoginBatchRunner: if proxy_msg: self._push_log("info", f"[{current}] {proxy_msg}") - # API代理模式下,验证代理池是否有可用代理;池空时等待恢复(最多 _PROXY_WAIT_MAX 秒) - # 注意:不预取代理传给 DouyuLogin,让 DouyuLogin 通过 proxy_manager 内部自治管理 - is_api_proxy = ( - self.proxy_config - and self.proxy_config.enabled - and not (self.proxy_config.http or self.proxy_config.https) - and self._shared_proxy_manager - ) - if is_api_proxy: - # 先验证代理池是否有可用代理(获取后立即归还,不占用) - test_proxy = self._shared_proxy_manager.get_proxy() - if test_proxy: - self._shared_proxy_manager.release_proxy(test_proxy) - self._push_log("info", f"[{current}] 代理池可用,由 DouyuLogin 内部管理代理获取与切换") - else: - # 代理池暂时耗尽,等待冷却代理恢复或新代理入池 - self._push_log("warning", f"[{current}] 代理池暂时耗尽,等待恢复...") - pool_available = False - for wait_sec in range(0, self._PROXY_WAIT_MAX, 10): - if self._stop.is_set(): - task.status = "error" - task.message = "任务已停止" - task.finished_at = datetime.now(timezone.utc) - worker_db.commit() - return - if self._sleep_or_stop(10): - task.status = "error" - task.message = "任务已停止" - task.finished_at = datetime.now(timezone.utc) - worker_db.commit() - return - test_proxy = self._shared_proxy_manager.get_proxy() - if test_proxy: - self._shared_proxy_manager.release_proxy(test_proxy) - pool_available = True - self._push_log("info", f"[{current}] 代理池恢复,由 DouyuLogin 内部管理代理") - break - remaining = self._PROXY_WAIT_MAX - wait_sec - 10 - self._push_log("info", f"[{current}] 代理池仍为空,继续等待... (剩余 {remaining}s)") - - if not pool_available: - task.status = "error" - task.message = f"代理池耗尽,等待 {self._PROXY_WAIT_MAX}s 后仍无可用代理" - task.finished_at = datetime.now(timezone.utc) - worker_db.commit() - self._push_log("error", f"[{current}] {acc_info['username']} 代理池耗尽,放弃") - return - # 静态代理启用但配置为空 → 不可用 - if self.proxy_config and self.proxy_config.enabled and not (self.proxy_config.http or self.proxy_config.https) and not self._shared_proxy_manager and not proxy_dict: + if self.proxy_config and self.proxy_config.enabled and not (self.proxy_config.http or self.proxy_config.https) and not self._shared_proxy_fetcher and not proxy_dict: task.status = "error" task.message = "代理不可用: 未配置代理" task.finished_at = datetime.now(timezone.utc) @@ -209,7 +152,7 @@ class LoginBatchRunner: max_proxy_retries=self.max_proxy_retries, max_login_retries=self.max_login_retries, max_total_time=self.max_total_time, - proxy_manager=self._shared_proxy_manager, + proxy_fetcher=self._shared_proxy_fetcher, stop_event=self._stop, ) result = loginer.login() @@ -333,24 +276,20 @@ class BatchRegistry: def __init__(self): self._batches: dict[str, dict] = {} - self._lock = threading.Lock() def register(self, batch_id: str, log_queue: asyncio.Queue, loop: asyncio.AbstractEventLoop, runner: LoginBatchRunner): - with self._lock: - self._batches[batch_id] = { - "log_queue": log_queue, - "loop": loop, - "runner": runner, - } + self._batches[batch_id] = { + "log_queue": log_queue, + "loop": loop, + "runner": runner, + } def get(self, batch_id: str): - with self._lock: - return self._batches.get(batch_id) + return self._batches.get(batch_id) def pop(self, batch_id: str): - with self._lock: - return self._batches.pop(batch_id, None) + return self._batches.pop(batch_id, None) # 模块级单例