File size: 5,904 Bytes
8a03d2c
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
"""
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 销毁