Spaces:
Running
Running
| """ | |
| Vcore AI 网络客户端 | |
| 负责处理底层的 HTTP 请求、连接池管理和重试逻辑。 | |
| 包含 Google Recaptcha Token 现抓现用逻辑。 | |
| """ | |
| import asyncio | |
| import contextlib | |
| import random | |
| import re | |
| from urllib.parse import parse_qs, urlparse | |
| from bs4 import BeautifulSoup | |
| from typing import Any, AsyncGenerator, Optional | |
| from curl_cffi import requests | |
| from curl_cffi.requests import Response | |
| from src.core.config import load_config | |
| from src.utils.logger import get_logger | |
| logger = get_logger(__name__) | |
| def _random_string(length: int) -> str: | |
| return "".join(random.choice("abcdefghijklmnopqrstuvwxyz0123456789") for _ in range(length)) | |
| class NetworkClient: | |
| """底层网络客户端""" | |
| def __init__(self): | |
| self.config = load_config() | |
| self.recaptcha_base_api = "https://www.google.com" | |
| self.browser_targets = ["chrome124", "chrome131"] | |
| logger.debug(f"NetworkClient 初始化完成") | |
| async def close(self): | |
| pass # 已改为即用即毁,无需全局清理 | |
| def _get_imp(self) -> str: | |
| return random.choice(self.browser_targets) | |
| async def fetch_recaptcha_token(self, session: requests.AsyncSession) -> Optional[str]: | |
| """获取一次 Google Recaptcha Token。 | |
| 重试预算由单节点协程统一管理;这里不再做固定 3 次内部重试。 | |
| """ | |
| random_cb = _random_string(10) | |
| anchor_url = f"{self.recaptcha_base_api}/recaptcha/enterprise/anchor?ar=1&k=6LdCjtspAAAAAMcV4TGdWLJqRTEk1TfpdLqEnKdj&co=aHR0cHM6Ly9jb25zb2xlLmNsb3VkLmdvb2dsZS5jb206NDQz&hl=zh-CN&v=ne1iDVwClkE7nKD3uA9Vqsvl&size=invisible&anchor-ms=20000&execute-ms=15000&cb={random_cb}" | |
| reload_url = f"{self.recaptcha_base_api}/recaptcha/enterprise/reload?k=6LdCjtspAAAAAMcV4TGdWLJqRTEk1TfpdLqEnKdj" | |
| try: | |
| anchor_response = await session.get(anchor_url, timeout=15) | |
| soup = BeautifulSoup(anchor_response.text, "html.parser") | |
| token_element = soup.find("input", {"id": "recaptcha-token"}) | |
| if token_element is None: | |
| logger.warning("anchor_html 未找到 token 元素") | |
| return None | |
| base_recaptcha_token = str(token_element.get("value")) | |
| parsed = urlparse(anchor_url) | |
| params = parse_qs(parsed.query) | |
| payload = { | |
| "v": params["v"][0], "reason": "q", "k": params["k"][0], | |
| "c": base_recaptcha_token, "co": params["co"][0], | |
| "hl": params["hl"][0], "size": "invisible", | |
| "vh": "6581054572", "chr": "", "bg": "", | |
| } | |
| headers = {"Content-Type": "application/x-www-form-urlencoded"} | |
| reload_response = await session.post( | |
| reload_url, data=payload, headers=headers, timeout=15 | |
| ) | |
| match = re.search(r'rresp","(.*?)"', reload_response.text) | |
| if not match: | |
| logger.warning("未找到 rresp") | |
| return None | |
| final_token = match.group(1) | |
| logger.debug("成功获取 Recaptcha Token") | |
| return final_token | |
| except Exception as e: | |
| with contextlib.suppress(Exception): | |
| setattr(session, "_vcore_proxy_last_recaptcha_error", str(e)) | |
| logger.debug(f"获取 recaptcha_token 失败: {e}") | |
| return None | |
| def create_session(self) -> requests.AsyncSession: | |
| """创建一个直连且带有随机伪装指纹的 Session。""" | |
| imp = self._get_imp() | |
| logger.debug(f"创建直连 Session (指纹: {imp})") | |
| return requests.AsyncSession(impersonate=imp, proxy=None) | |
| def create_session_with_proxy(self, proxy_url: Optional[str]) -> requests.AsyncSession: | |
| """创建使用指定代理的独立 Session,用于并行节点请求。""" | |
| imp = self._get_imp() | |
| logger.debug(f"创建节点专用 Session (指纹: {imp}, 代理: {proxy_url or '直连'})") | |
| return requests.AsyncSession(impersonate=imp, proxy=proxy_url) | |
| async def post_request(self, session: requests.AsyncSession, url: str, headers: dict[str, str], json_data: dict[str, Any]) -> Response: | |
| """发送非流式 POST 请求 (复用 Session)""" | |
| try: | |
| return await session.post(url=url, headers=headers, json=json_data, timeout=180.0) | |
| except asyncio.CancelledError: | |
| logger.debug("非流式网络请求被取消,立即关闭上游 Session") | |
| try: | |
| await session.close() | |
| except Exception as close_error: | |
| logger.debug(f"取消非流式请求时关闭 Session 失败: {close_error}") | |
| raise | |
| except Exception as e: | |
| logger.debug(f"非流式网络请求异常: {e}") | |
| raise | |
| async def stream_request( | |
| self, | |
| session: requests.AsyncSession, | |
| method: str, | |
| url: str, | |
| headers: dict[str, str], | |
| json_data: dict[str, Any], | |
| ) -> AsyncGenerator[Response, None]: | |
| """发送流式请求 (复用 Session),响应首包后由调用方直接转发底层结果。""" | |
| try: | |
| response = await session.request( | |
| method=method, url=url, headers=headers, json=json_data, timeout=180.0, stream=True | |
| ) | |
| yield response | |
| except asyncio.CancelledError: | |
| logger.debug("流式网络请求被取消,立即关闭上游 Session") | |
| try: | |
| await session.close() | |
| except Exception as close_error: | |
| logger.debug(f"取消流式请求时关闭 Session 失败: {close_error}") | |
| raise | |
| except Exception as e: | |
| logger.debug(f"网络请求异常: {e}") | |
| raise | |
| # 注意:这里不再主动关闭 response,由调用方在读取完流后处理,或者随 session 销毁 | |