Files
douyin/douyin-im-collector/douyin_collector_server.py
2026-07-17 09:24:47 +08:00

398 lines
14 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""抖音一键采集器 · 本地图形化页面服务(stdlib,无额外依赖)"""
from __future__ import annotations
import json
import subprocess
import sys
import threading
import time
import uuid
from http.server import BaseHTTPRequestHandler, HTTPServer
from pathlib import Path
from typing import Any
from urllib.parse import urlparse
from app_paths import ui_html_path
UI_FILE = ui_html_path()
DEFAULT_PORT = 8765
class CollectorState:
def __init__(self) -> None:
self._lock = threading.Lock()
self.reset()
def reset(self) -> None:
with self._lock:
self.phase = "waiting"
self.message = "正在启动…"
self.result = None
self.checks: list[dict[str, str]] = []
self.result_text = ""
self.my_uid = ""
self.ok_core = False
self.summary = ""
self.run_id = uuid.uuid4().hex[:8]
self.started_at = int(time.time())
self.port = DEFAULT_PORT
self.step = 1
self.logs: list[str] = []
def append_log(self, message: str) -> None:
line = (message or "").strip()
if not line:
return
with self._lock:
self.logs.append(line)
if len(self.logs) > 80:
self.logs = self.logs[-80:]
def set_step(self, step: int) -> None:
with self._lock:
self.step = max(1, min(4, step))
def set_waiting(self, message: str = "") -> None:
with self._lock:
self.phase = "waiting"
if message:
self.message = message
if self.step < 2:
self.step = 2
def set_collecting(self, message: str = "") -> None:
with self._lock:
self.phase = "collecting"
self.message = message or "正在采集 Cookie / 签名 / WebSocket…"
self.step = 3
def set_done(self, result: dict, checks: list[dict[str, str]], ok_core: bool, summary: str) -> None:
with self._lock:
self.phase = "done"
self.result = result
self.checks = checks
self.result_text = json.dumps(result, ensure_ascii=False, indent=2)
self.my_uid = str(result.get("my_uid") or "")
self.ok_core = ok_core
self.summary = summary
self.message = "采集完成"
self.step = 4
def set_error(self, message: str) -> None:
with self._lock:
self.phase = "error"
self.message = message
def snapshot(self) -> dict[str, Any]:
with self._lock:
return {
"phase": self.phase,
"message": self.message,
"checks": self.checks,
"result_text": self.result_text,
"my_uid": self.my_uid,
"ok_core": self.ok_core,
"summary": self.summary,
"run_id": self.run_id,
"started_at": self.started_at,
"port": self.port,
"step": self.step,
"logs": list(self.logs),
}
def append_log(message: str) -> None:
STATE.append_log(message)
STATE = CollectorState()
_SERVER: HTTPServer | None = None
_start_event: threading.Event | None = None
_shutdown_event = threading.Event()
def bind_start_event(event: threading.Event) -> None:
"""采集主流程注册:网页「开始采集」或自动就绪时触发。"""
global _start_event
_start_event = event
def request_start_collect() -> bool:
if _start_event is None:
return False
_start_event.set()
return True
def is_start_requested() -> bool:
return bool(_start_event and _start_event.is_set())
def request_shutdown() -> None:
_shutdown_event.set()
srv = _SERVER
if srv:
threading.Thread(target=srv.shutdown, daemon=True).start()
def wait_for_shutdown() -> None:
_shutdown_event.wait()
def _free_port(port: int) -> None:
"""关闭占用端口的旧采集器进程(避免网页显示上一次失败的错误)。"""
if sys.platform != "win32":
return
try:
subprocess.run(
[
"powershell",
"-NoProfile",
"-ExecutionPolicy",
"Bypass",
"-Command",
(
f"Get-NetTCPConnection -LocalPort {port} -State Listen "
f"-ErrorAction SilentlyContinue | ForEach-Object {{ "
f"Stop-Process -Id $_.OwningProcess -Force -ErrorAction SilentlyContinue }}"
),
],
capture_output=True,
timeout=10,
creationflags=getattr(subprocess, "CREATE_NO_WINDOW", 0),
)
except Exception:
pass
def _parse_web_protect(raw: str) -> dict:
if not raw:
return {}
try:
obj = json.loads(str(raw))
if isinstance(obj, str):
obj = json.loads(obj)
data = obj.get("data") if isinstance(obj, dict) else None
if isinstance(data, str):
return json.loads(data)
if isinstance(data, dict):
return data
return obj if isinstance(obj, dict) else {}
except Exception:
return {}
def _pick_ls(result: dict, exact: str, contains: tuple[str, ...] = ()) -> str:
for origin in result.get("origins") or []:
for entry in origin.get("localStorage") or []:
name = entry.get("name") or ""
val = entry.get("value") or ""
if name == exact and val:
return str(val)
lower = name.lower()
if val and contains and any(c in lower for c in contains):
return str(val)
return ""
def _extract_my_uid(result: dict) -> str:
uid = result.get("my_uid")
if uid and str(uid).isdigit():
return str(uid)
ws_url = str(result.get("frontier_ws_url") or "")
if "aid=2906" in ws_url:
try:
from urllib.parse import parse_qs, urlparse
dev = parse_qs(urlparse(ws_url).query).get("device_id", [""])[0]
if dev and str(dev).isdigit():
return str(dev)
except Exception:
pass
tea = _pick_ls(result, "tea_cache_tokens", ("tea_cache_tokens",))
if tea:
try:
uid = json.loads(tea).get("user_unique_id")
if uid and str(uid).isdigit():
return str(uid)
except Exception:
pass
return ""
def normalize_result(result: dict) -> dict:
try:
from credential_normalize import normalize_storage_state_for_im
return normalize_storage_state_for_im(result)
except Exception:
return result
def build_checks(result: dict) -> tuple[list[dict[str, str]], bool, str]:
result = normalize_result(result)
cookies = result.get("cookies") or []
names = {str(c.get("name", "")).lower() for c in cookies}
val_protect = _pick_ls(result, "security-sdk/s_sdk_sign_data_key/web_protect", ("web_protect",))
val_keys = _pick_ls(result, "security-sdk/s_sdk_crypt_sdk", ("crypt_sdk",))
my_uid = _extract_my_uid(result)
wp = _parse_web_protect(val_protect)
has_session = "sessionid" in names or "sessionid_ss" in names
has_protect = bool(val_protect)
has_keys = bool(val_keys)
has_uid = bool(my_uid)
has_wp_full = bool(wp.get("ticket") and wp.get("ts_sign") and (wp.get("client_cert") or wp.get("sdk_cert")))
has_ws = bool(result.get("frontier_ws_url"))
has_sdk_cert = bool(result.get("sdk_cert"))
token = ""
ws_url = str(result.get("frontier_ws_url") or "")
if ws_url:
try:
from urllib.parse import parse_qs, urlparse
token = parse_qs(urlparse(ws_url).query).get("token", [""])[0]
except Exception:
pass
ws_real = has_ws and (has_sdk_cert or len(token) >= 40)
def chk(ok: bool, ok_t: str, fail_t: str, critical: bool = True) -> dict[str, str]:
if ok:
return {"s": "ok", "t": "已包含 " + ok_t}
return {"s": "err" if critical else "warn", "t": fail_t}
checks = [
chk(has_session, "sessionid(登录态)", "缺 sessionid:请重新登录并在私信页停留后采集", True),
chk(has_uid, f"抖音 UID{my_uid}" if my_uid else "抖音 UID", "未提取到数字 UID", False),
chk(has_protect, "web_protect 签名", "缺 web_protect:请在私信页停留 510 秒后重采", True),
]
if has_protect:
checks.append(chk(has_wp_full, "web_protect 完整字段", "web_protect 不完整(缺 ticket/ts_sign/cert", True))
checks.append(chk(has_keys, "crypt_sdk 密钥", "缺 crypt_sdk:请确认在已登录页面采集", True))
if has_ws and has_sdk_cert:
aid = "2906 创作者私信" if "aid=2906" in ws_url else "真实抓包"
checks.append({"s": "ok", "t": f"frontier WebSocket{aid},含 sdk_cert"})
elif ws_real:
checks.append({"s": "ok", "t": "frontier WebSocket(真实抓包)"})
elif has_ws:
checks.append({"s": "warn", "t": "WS 无效(token 异常),请重新在私信页采集"})
elif has_session:
checks.append({"s": "warn", "t": "未捕获 WS:请在私信页等消息列表加载后重采"})
else:
checks.append({"s": "warn", "t": "未捕获 WS,请先完成登录再采集"})
ok_core = has_session and has_protect and has_keys and has_wp_full and has_uid
if ok_core:
summary = f"共 {len(cookies)} 条 Cookie,可直接复制导入系统。"
elif has_session and has_protect and has_keys:
summary = "核心凭证已有,建议核对红色项后导入;若 IM 异常请重新采集。"
else:
summary = "关键凭证缺失,请按提示重新登录采集,不要直接导入。"
return checks, ok_core, summary
class _Handler(BaseHTTPRequestHandler):
def log_message(self, format: str, *args) -> None: # noqa: A003
pass
def _send(self, code: int, body: bytes, content_type: str, *, cors: bool = False) -> None:
self.send_response(code)
self.send_header("Content-Type", content_type)
self.send_header("Cache-Control", "no-store")
if cors:
self.send_header("Access-Control-Allow-Origin", "*")
self.send_header("Access-Control-Allow-Methods", "GET, POST, OPTIONS")
self.send_header("Access-Control-Allow-Headers", "Content-Type")
self.send_header("Content-Length", str(len(body)))
self.end_headers()
self.wfile.write(body)
def do_OPTIONS(self) -> None: # noqa: N802
self.send_response(204)
self.send_header("Access-Control-Allow-Origin", "*")
self.send_header("Access-Control-Allow-Methods", "GET, POST, OPTIONS")
self.send_header("Access-Control-Allow-Headers", "Content-Type")
self.end_headers()
def do_GET(self) -> None: # noqa: N802
path = urlparse(self.path).path
if path in ("/", "/index.html"):
if UI_FILE.is_file():
self._send(200, UI_FILE.read_bytes(), "text/html; charset=utf-8")
else:
self._send(404, b"UI not found", "text/plain")
return
if path == "/api/status":
payload = json.dumps(STATE.snapshot(), ensure_ascii=False).encode("utf-8")
self._send(200, payload, "application/json; charset=utf-8", cors=True)
return
if path == "/api/ping":
self._send(200, b"ok", "text/plain; charset=utf-8", cors=True)
return
self._send(404, b"Not Found", "text/plain")
def do_POST(self) -> None: # noqa: N802
path = urlparse(self.path).path
if path == "/api/start":
ok = request_start_collect()
if ok:
STATE.set_collecting("收到开始信号,正在采集…")
append_log("用户点击:立即采集")
body = json.dumps({"ok": ok}, ensure_ascii=False).encode("utf-8")
self._send(200 if ok else 409, body, "application/json; charset=utf-8", cors=True)
return
if path == "/api/quit":
append_log("用户退出采集器")
request_shutdown()
body = json.dumps({"ok": True}, ensure_ascii=False).encode("utf-8")
self._send(200, body, "application/json; charset=utf-8", cors=True)
return
self._send(404, b"Not Found", "text/plain")
def start_server(port: int = DEFAULT_PORT) -> tuple[HTTPServer, str]:
global _SERVER
STATE.reset()
STATE.set_step(1)
append_log("正在启动本地界面…")
for attempt in range(10):
p = port + attempt
_free_port(p)
try:
srv = HTTPServer(("127.0.0.1", p), _Handler)
_SERVER = srv
STATE.port = p
url = f"http://127.0.0.1:{p}/"
thread = threading.Thread(target=srv.serve_forever, daemon=True)
thread.start()
return srv, url
except OSError:
continue
raise RuntimeError("无法启动本地页面服务(端口被占用)")
def open_browser(url: str) -> None:
try:
if sys.platform == "win32":
subprocess.Popen(
["cmd", "/c", "start", "", url],
creationflags=getattr(subprocess, "CREATE_NO_WINDOW", 0),
close_fds=True,
)
return
except Exception:
pass
try:
import webbrowser
webbrowser.open(url)
except Exception:
pass
def publish_result(result: dict) -> None:
normalized = normalize_result(result)
checks, ok_core, summary = build_checks(normalized)
STATE.set_done(normalized, checks, ok_core, summary)
def publish_error(message: str) -> None:
STATE.set_error(message)