# -*- coding: utf-8 -*-
"""DouParse 自建代理 agent（在你自己的电脑上运行，把住宅 IP 借给服务器抓抖音）。

原理：家庭电脑在 NAT 后面，服务器主动连不进来，所以由 agent **主动** WebSocket
连出去到服务器（穿透 NAT）。服务器直连 IP 被抖音封禁时，把"抓分享页"任务顺着
这条隧道派给 agent，agent 用本地住宅 IP 抓取并回传 HTML。

协议与浏览器扩展中继完全一致（/ws/relay）：
  服务器 → {type:"task", id, url}
  agent  → {type:"result", id, ok, html, finalUrl}
  agent  → {type:"ping"}  服务器 → {type:"pong"}

登录态自动同步：连上后服务器会下发身份池的登录 Cookie（config/
pong 消息携带，wss 加密），agent 自动带登录态抓取，成功率与服务器
对齐；无需手动配置。

运行（在本机 backend 目录，用现有 venv，无需额外安装）：
    .venv\\Scripts\\python.exe proxy_agent.py
可选环境变量：
    DP_WS=wss://www.qiaiwan.top/ws/relay   （默认即此）
    DP_COOKIE=<一条登录 Cookie>             （可选：带上登录态抓取成功率更高，
                                              与服务器身份池同理；不配也能跑）

单实例保护：同机只允许跑一个 agent（同 IP 多开只会白白浪费自家宽带配额），
重复启动会立即退出；配合开机自启动脚本可放心常驻后台。
日志同时写入本目录 proxy_agent.log（pythonw 无窗口后台运行时看这里）。

注意：agent 会消耗你家宽带的抖音请求配额，量大会封你家 IP，请按需启停。
"""
from __future__ import annotations

import asyncio
import json
import logging
import os
import random
import re
import socket
import sys
import time

import httpx
import websockets

# 版本号（引导脚本安装时打印，方便确认机器上跑的是不是新版；每次改动递增）
AGENT_VERSION = "2026.08.13.2"

WS_URL = os.environ.get("DP_WS", "wss://www.qiaiwan.top/ws/relay?kind=agent")
# 来源标识：py=自建代理进程（后台与扩展内代理 ext 区分统计）
if "src=" not in WS_URL:
    WS_URL += ("&" if "?" in WS_URL else "?") + "src=py"

# 被服务器顶替/隔离后的冷静期（同 IP 多客户端抢名额时避免踢-重连死循环）
hold_until = 0.0

_BASE_DIR = os.path.dirname(os.path.abspath(__file__))
_LOG_FILE = os.path.join(_BASE_DIR, "proxy_agent.log")

# 单实例锁：绑定固定本地端口，占得住就是第一个实例；开机自启/手动误点多开都安全。
_LOCK_PORT = 57821
_lock_sock: socket.socket | None = None


def _acquire_single_instance() -> bool:
    global _lock_sock
    try:
        s = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
        s.bind(("127.0.0.1", _LOCK_PORT))
        _lock_sock = s  # 持有引用直到进程退出，端口即锁
        return True
    except OSError:
        return False


# 日志落盘 + 有控制台时同时输出（pythonw 后台无窗口时 sys.stderr 为 None）
_handlers: list[logging.Handler] = [logging.FileHandler(_LOG_FILE, encoding="utf-8")]
if sys.stderr is not None:
    _handlers.append(logging.StreamHandler())
logging.basicConfig(level=logging.INFO, format="%(asctime)s %(levelname)s %(message)s", handlers=_handlers)
logging.getLogger("httpx").setLevel(logging.WARNING)  # 只留任务级日志，避免每条 HTTP 都刷屏

MOBILE_UAS = [
    "Mozilla/5.0 (iPhone; CPU iPhone OS 17_0 like Mac OS X) AppleWebKit/605.1.15 (KHTML, like Gecko) Version/17.0 Mobile/15E148 Safari/604.1",
    "Mozilla/5.0 (iPhone; CPU iPhone OS 17_2 like Mac OS X) AppleWebKit/605.1.15 (KHTML, like Gecko) Version/17.2 Mobile/15E148 Safari/604.1",
    "Mozilla/5.0 (iPhone; CPU iPhone OS 17_4 like Mac OS X) AppleWebKit/605.1.15 (KHTML, like Gecko) Version/17.4 Mobile/15E148 Safari/604.1",
    "Mozilla/5.0 (iPhone; CPU iPhone OS 17_5_1 like Mac OS X) AppleWebKit/605.1.15 (KHTML, like Gecko) Version/17.5 Mobile/15E148 Safari/604.1",
    "Mozilla/5.0 (iPhone; CPU iPhone OS 18_0 like Mac OS X) AppleWebKit/605.1.15 (KHTML, like Gecko) Version/18.0 Mobile/15E148 Safari/604.1",
]
# 兼容旧配置引用
MOBILE_UA = MOBILE_UAS[0]

log = logging.getLogger("proxy_agent")

# 风控重定向目标 URL 里的关键段（命中即视为验证墙）；HTML 内容里的
# captcha 等字样不可作依据——正常页面 JS 里也有这些字符串。
_VERIFY_URL_MARKERS = ("verify", "captcha", "/security/", "secsdk", "verify-center")

# 住宅 IP 经不起多路并发猛打（并发越高越容易触发抖音风控，把整个出口打废）：
# 限制同时抓取数，队列里的任务稍等而非拒绝（服务器派发超时 8s，3 路并发内够用）。
_FETCH_SEM = asyncio.Semaphore(3)

# 可选登录 Cookie（DP_COOKIE）：匿名 ttwid 容易被软风控掉包，带上登录态
# 后抓取成功率显著提升（与服务器身份池同理，2026-08-13 优化）。
# 注：服务器现在会通过 WS 自动下发登录 Cookie（优先级更高），
# DP_COOKIE 仅作为断网重连前的本地兜底。
_LOGIN_COOKIE = (os.environ.get("DP_COOKIE") or "").strip()

# 服务器下发的登录 Cookie（连接后 config 下发 + 每次心跳 pong 刷新）
_server_cookie = ""


def _note_server_cookie(msg: dict) -> None:
    """应用服务器下发的登录 Cookie（加密隧道内传输，仅限自建代理节点）。"""
    global _server_cookie
    ck = str(msg.get("cookie") or "").strip()
    if not ck:
        return
    if ck != _server_cookie:
        _server_cookie = ck
        log.info("server-pushed login cookie applied (len=%d)", len(ck))
    if _client is not None and not _client.is_closed:
        _client.headers["Cookie"] = ck

# ttwid 有效期有限：常驻进程启动时暖一次用一辈子的话，过期后所有抓取
# 都带着无效身份 → 每小时重新暖一次首页刷新 ttwid（2026-08-13 优化）。
_TTWID_REFRESH_SECONDS = 3600.0

# 持久客户端：ttwid 等 Cookie 跨任务复用，不用每次先暖一次首页
# （旧版每任务额外多打一次抖音，既慢又双倍消耗住宅 IP 配额）。
_client: httpx.AsyncClient | None = None
_client_warmed = False
_last_warm_ts = 0.0


def _apply_cookie(c: httpx.AsyncClient) -> None:
    """优先用服务器下发的登录 Cookie，其次本地 DP_COOKIE（都没有则匿名 ttwid）。"""
    ck = _server_cookie or _LOGIN_COOKIE
    if ck:
        c.headers["Cookie"] = ck


async def _get_client() -> httpx.AsyncClient:
    global _client, _client_warmed, _last_warm_ts
    if _client is None or _client.is_closed:
        _client = httpx.AsyncClient(
            headers={
                "User-Agent": random.choice(MOBILE_UAS),
                "Accept": "text/html,application/xhtml+xml,application/xml;q=0.9,*/*;q=0.8",
                "Accept-Language": "zh-CN,zh;q=0.9",
                "Referer": "https://www.douyin.com/",
                # 补齐真实浏览器的顶级导航元数据，与服务器 _client() 一致，
                # 降低指纹异常被软风控的概率
                "Sec-Fetch-Dest": "document",
                "Sec-Fetch-Mode": "navigate",
                "Sec-Fetch-Site": "none",
                "Sec-Fetch-User": "?1",
                "Upgrade-Insecure-Requests": "1",
            },
            timeout=12.0,
            follow_redirects=True,
        )
        _apply_cookie(_client)
        _client_warmed = False
    now = time.time()
    if not _client_warmed or now - _last_warm_ts > _TTWID_REFRESH_SECONDS:
        # 首次/定期重暖首页刷新 ttwid；顺手换个 UA 避免长期单一指纹
        _client.headers["User-Agent"] = random.choice(MOBILE_UAS)
        try:
            await _client.get("https://www.iesdouyin.com/")
        except Exception:  # noqa: BLE001
            pass
        _client_warmed = True
        _last_warm_ts = now
    return _client


def _candidate_urls(url: str) -> list[str]:
    """分享页候选 URL：原 URL 优先，失败后轮换移动域名重抓
    （与服务器 _SHARE_HOSTS 轮换同思路：单个域名被软风控时换域名
    往往就能拿到真数据，2026-08-13 优化）。短链等非分享页 URL
    只有自身一个候选，行为与旧版一致。"""
    urls = [url]
    m = re.search(r"/share/(?:video|slides)/(\d+)", url)
    if m:
        kind = "slides" if "/slides/" in url else "video"
        for h in ("m.douyin.com", "www.iesdouyin.com", "www.douyin.com"):
            cand = f"https://{h}/share/{kind}/{m.group(1)}/"
            if cand != url:
                urls.append(cand)
    return urls


async def fetch_share(url: str) -> dict:
    """用本地住宅 IP 抓抖音分享页，返回 {ok, html, finalUrl}。

    候选域名轮换：第一个 URL 拿不到作品数据时自动试其它移动域名
    （旧版只在空壳页时补试一次 m.douyin，现对无 SSR/软风控同样生效）。

    回传前自检：验证墙/无作品数据的软风控页面如实报 ok=False（错误
    文案带“风控”关键字），服务器会据此隔离本节点让住宅 IP 休息，
    避免拿垃圾页面去浪费其它候选域名。注意：正常页面的 JS 里也含
    captcha 等字样，所以 HTML 关键字不作判断依据，只看最终 URL 与
    页面里有没有作品数据。
    """
    try:
        c = await _get_client()
    except Exception as e:  # noqa: BLE001
        return {"ok": False, "error": str(e)[:80]}
    r = None
    last_err = ""
    async with _FETCH_SEM:
        for u in _candidate_urls(url):
            try:
                rr = await c.get(u)
            except Exception as e:  # noqa: BLE001
                last_err = str(e)[:80]
                continue
            # 拿到作品数据就用它；否则先留作兜底，继续试下一个候选域名
            if rr.text and ('"item_list"' in rr.text or '"filter_list"' in rr.text):
                r = rr
                break
            r = rr
    if r is None:
        return {"ok": False, "error": last_err or "全部候选域名抓取失败"}
    html = r.text
    final_url = str(r.url)
    if not html or len(html) > 4 * 1024 * 1024:
        return {"ok": False, "error": "页面过大或为空"}
    low_url = final_url.lower()
    if any(m in low_url for m in _VERIFY_URL_MARKERS):
        return {"ok": False, "error": f"本出口 IP 命中抖音风控验证页 final={final_url[:80]}"}
    # 软风控：200 但页面里没有作品数据（_ROUTER_DATA/RENDER_DATA 都没有，
    # 或有 SSR 但既无 item_list 也无 filter_list）——回传也没用，如实上报。
    has_ssr = ("_ROUTER_DATA" in html) or ("RENDER_DATA" in html)
    has_item_key = ('"item_list"' in html) or ('"filter_list"' in html)
    if not (has_ssr and has_item_key):
        # 有 SSR 但无作品键 = 空壳页（非风控，不报“风控”字样避免节点被冷却）；
        # 连 SSR 都没有 = 验证墙/掉包页，报风控让服务器冷却本节点。
        if has_ssr:
            return {"ok": False, "error": f"桌面空壳页无作品数据（len={len(html)}）"}
        return {"ok": False, "error": f"页面不含作品数据（疑似风控，len={len(html)}）"}
    return {"ok": True, "html": html, "finalUrl": final_url}


async def handle_task(ws, msg: dict) -> None:
    tid = str(msg.get("id") or "")
    url = str(msg.get("url") or "")
    res = await fetch_share(url)
    try:
        # error 必须随结果回传：服务器靠它区分“风控失败”与普通过错，
        # 进而决定让节点冷却休息还是继续派单（旧版漏发此字段，服务器
        # 永远拿到空 err，风控冷却分支从未生效）。
        await ws.send(json.dumps(
            {"type": "result", "id": tid, "ok": bool(res.get("ok")),
             "html": res.get("html", ""), "finalUrl": res.get("finalUrl", ""),
             "error": res.get("error", "")}
        ))
        if res.get("ok"):
            log.info("task %s -> ok len=%d final=%s", tid[:6], len(res.get("html") or ""), res.get("finalUrl") or "")
        else:
            log.warning("task %s -> FAIL %s url=%s", tid[:6], res.get("error"), url[:80])
    except Exception:  # noqa: BLE001
        pass


async def pinger(ws, rx: dict) -> None:
    while True:
        await asyncio.sleep(25)
        # pong 看门狗：服务器对每个 ping 都回 pong，超 90s 没收到任何消息
        # 说明连接已死（如服务器重启后客户端无感知）→ 主动断开触发重连
        if time.time() - rx["t"] > 90:
            log.warning("server silent >90s, force reconnect")
            try:
                await ws.close()
            except Exception:  # noqa: BLE001
                pass
            return
        try:
            await ws.send(json.dumps({"type": "ping"}))
        except Exception:  # noqa: BLE001
            return


async def run_once() -> None:
    global hold_until
    log.info("connect %s", WS_URL)
    async with websockets.connect(WS_URL, max_size=8 * 1024 * 1024) as ws:
        welcome = json.loads(await ws.recv())
        if welcome.get("type") == "reject":
            log.warning("rejected by server, wait 60s")
            await asyncio.sleep(60)
            return
        log.info("connected as node %s", welcome.get("node"))
        rx = {"t": time.time()}  # 最近一次收到服务器消息的时间（看门狗用）
        ping_task = asyncio.create_task(pinger(ws, rx))
        try:
            async for raw in ws:
                rx["t"] = time.time()
                try:
                    msg = json.loads(raw)
                except Exception:
                    continue
                if msg.get("type") == "reject":
                    # 被同 IP 新连接顶替 / 隔离：长退避 10 分钟，
                    # 否则“踢→3s重连→再踢”死循环，节点永远活不了
                    log.warning("server reject (replaced/quarantined), hold 10min")
                    hold_until = time.time() + 600
                    return
                if msg.get("type") in ("config", "pong"):
                    # 服务器下发/刷新登录 Cookie：带上登录态抓取，成功率对齐服务器
                    _note_server_cookie(msg)
                    continue
                if msg.get("type") == "task":
                    asyncio.create_task(handle_task(ws, msg))
        finally:
            ping_task.cancel()


async def main() -> None:
    backoff = 3
    while True:
        wait = hold_until - time.time()
        if wait > 0:
            await asyncio.sleep(wait)
            continue
        try:
            await run_once()
            backoff = 3
        except Exception as e:  # noqa: BLE001
            log.warning("disconnected: %s, retry in %ss", e, backoff)
            await asyncio.sleep(backoff + random.random())
            backoff = min(backoff * 2, 30)


if __name__ == "__main__":
    if not _acquire_single_instance():
        log.warning("another proxy_agent is already running on this machine, exit")
        sys.exit(0)
    try:
        asyncio.run(main())
    except KeyboardInterrupt:
        log.info("stopped")
