diff --git a/SCRAPE_MODULE_DEV_GUIDE.md b/SCRAPE_MODULE_DEV_GUIDE.md new file mode 100644 index 0000000..9565376 --- /dev/null +++ b/SCRAPE_MODULE_DEV_GUIDE.md @@ -0,0 +1,874 @@ +# 内容抓取模块 · 开发技术说明 + +> 面向接手开发的团队 · 2026-10-10 +> 全部内容基于**逐行读源码**整理,不是推测 +> 范围:**只写内容抓取模块**,不涉及其他业务 + +--- + +## 一、这个模块是干什么的 + +从主流内容平台(抖音 / 小红书 / B站 / 知乎 / 任意网页)**采集内容数据**: +- 视频/笔记的元数据(标题、作者、发布时间、正文) +- 互动数据(点赞、评论、分享、播放) +- 评论列表 +- 作者主页的全部作品列表 + +采到的数据落进本地 SQLite,供后续分析/运营使用。 + +--- + +## 二、整体架构(重要:三层降级是核心) + +``` +┌──────────────────────────────────────────────────────────┐ +│ HTTP 层 routes/scrape.py(22 个端点) │ +└───────────────────────┬──────────────────────────────────┘ + ▼ +┌──────────────────────────────────────────────────────────┐ +│ 引擎层 services/scrape_engine.py │ +│ · URL 解析 → 标准化目标 │ +│ · 按平台选适配器 │ +│ · 同步 / 异步调度 │ +│ · 结果落库(去重) │ +└───────────────────────┬──────────────────────────────────┘ + ▼ +┌──────────────────────────────────────────────────────────┐ +│ 适配器层 services/adapters/*.py(每个平台一个) │ +│ 基类 ScrapeAdapter 定义统一接口 + 工具降级 │ +└───────────────────────┬──────────────────────────────────┘ + ▼ +┌──────────────────────────────────────────────────────────┐ +│ 工具层(★ 三层降级,这是本模块的核心设计) │ +│ Level 1 OpenCLI 外部 Node CLI(主路径) │ +│ Level 2 agent-browser Playwright + 真实 Chrome │ +│ Level 3 web_crawler 通用网页兜底 │ +└───────────────────────┬──────────────────────────────────┘ + ▼ +┌──────────────────────────────────────────────────────────┐ +│ 存储层 services/scrape_db.py(SQLite,4 张表) │ +└──────────────────────────────────────────────────────────┘ +``` + +### 为什么这么设计 + +抓取的最大风险是**单一方式失效**:目标平台改版、接口封禁、登录态过期。 +所以**同一份数据有三条获取路径**,第一条失败自动降级到第二条, +**上层完全不感知**(对 engine 来说只是"拿到数据了")。 + +--- + +## 三、工具层详解(最关键的一层) + +### 3.1 Level 1 — OpenCLI + +**它是什么**:一个**第三方 Node.js CLI 工具**,包名 `@jackwener/opencli`。 + +``` +实际安装位置(本机实测): + ~/.workbuddy/binaries/node/versions/22.22.2/bin/opencli + → 软链到 ../lib/node_modules/@jackwener/opencli/dist/src/main.js +``` + +**怎么调用**(`services/adapters/__init__.py` 的 `_run_opencli`): + +```python +OPENCLI = os.environ.get("OPENCLI_PATH", + str(Path.home() / ".workbuddy" / "binaries" / "node" / + "versions" / "22.22.2" / "bin" / "opencli")) + +async def _run_opencli(self, args: list, timeout: int = 60): + cmd = [self.OPENCLI] + args + proc = await asyncio.create_subprocess_exec( + *cmd, stdout=PIPE, stderr=PIPE) + stdout, stderr = await asyncio.wait_for(proc.communicate(), timeout=timeout) + ... + return self._parse_output(stdout.decode().strip()) +``` + +**调用示例**(抖音,`douyin_scrape.py`): +```bash +opencli douyin user-videos --limit 20 --with_comments true -f json +opencli douyin stats -f json +``` + +**⛔ 移植注意**:这个二进制**不在仓库里**,是外部依赖。移植时必须: +- 要么在目标机装 `npm i -g @jackwener/opencli` +- 要么改 `OPENCLI_PATH` 环境变量指向它的位置 + +### 3.2 Level 2 — agent-browser("套用真实浏览器"的做法) + +**这是你问的重点。设计原则写在 `browser_helpers.py` 文件头**: + +``` +⛔ 绝不使用 Camoufox(养号专用,Firefox 内核 + 特殊指纹) +✅ 使用 Playwright 启动【真实 Chrome】(Chromium 内核,正常指纹) +``` + +**具体怎么"套真实浏览器"**(`browser_helpers.py` 的 `_get_browser`): + +```python +_CHROME_PATHS = [ + "/Applications/Google Chrome.app/Contents/MacOS/Google Chrome", # ← 系统真 Chrome + "/Applications/Chromium.app/Contents/MacOS/Chromium", +] +_CHROME_PATH = None +for p in _CHROME_PATHS: + if Path(p).exists(): + _CHROME_PATH = p # 自动探测,找到就用系统已装的 Chrome + break + +launch_kwargs = { + "headless": headless, + "args": [ + "--disable-blink-features=AutomationControlled", # ★ 反检测关键 + "--no-sandbox", + "--disable-dev-shm-usage", + "--disable-gpu", + "--window-size=1280,720", + ], +} +if _CHROME_PATH: + launch_kwargs["executable_path"] = _CHROME_PATH # ★ 用系统 Chrome,不用 Playwright 自带 +``` + +**四个关键设计点**: + +| 点 | 做法 | 为什么 | +|---|---|---| +| **用什么内核** | 系统真实 Chrome(`executable_path` 指定) | Playwright 自带 Chromium 有明显特征;真实 Chrome 是正常用户指纹 | +| **怎么隐藏自动化** | `--disable-blink-features=AutomationControlled` | 这是最常被检测的自动化标志位 | +| **实例管理** | 模块级单例 `_browser` + `asyncio.Lock` | 避免每次请求都启动浏览器(启动 ~1-2 秒) | +| **会话隔离** | 每次 `browser.new_context()` | 每个任务独立 cookie 环境,互不污染 | + +**页面加载策略**(`page_evaluate`): +```python +await page.goto(url, wait_until="domcontentloaded", timeout=timeout) +await page.wait_for_load_state("networkidle", timeout=timeout) # 等动态渲染 +await asyncio.sleep(1) # 再等 1 秒保险 +result = await page.evaluate(js_code) # 执行 JS 提取 +``` + +**UA 伪装**(每次 context 都设置): +```python +context = await browser.new_context( + user_agent=("Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) " + "AppleWebKit/537.36 (KHTML, like Gecko) " + "Chrome/125.0.0.0 Safari/537.36"), + viewport={"width": 1280, "height": 720}, + locale="zh-CN", # 中文环境,符合目标用户画像 +) +``` + +**两个公开函数**: +```python +page_evaluate(url, js_code, timeout, headless) # 打开页面执行 JS,返回 dict +page_extract(url, selectors, timeout, headless) # 按 CSS 选择器提取文本 +``` + +**Profile 持久化**(可选,环境变量控制): +```python +_USER_DATA_DIR = os.environ.get( + "SCRAPE_CHROME_USER_DATA", + str(Path.home() / "workbuddy-agent-os" / "agent-local" / + "runtime" / "scrape_chrome_profile")) +``` +→ 想保持登录态就指定这个目录;不指定就用临时 context。 + +### 3.3 Level 3 — web_crawler + +通用网页抓取兜底(`web_scrape.py`),用于非主流平台的页面。 + +### 3.4 降级怎么触发(`ScrapeAdapter._try_tools`) + +```python +async def _try_tools(self, tool_level: int, funcs: list) -> tuple: + """funcs: [(工具名, 可调用对象), ...],按顺序尝试""" + tools = [f for f in funcs[:tool_level]] # 按 level 截断 + for name, fn in tools: + try: + result = await fn() + if result is not None: + return True, result, name # ★ 成功即返回,不再降级 + except Exception as e: + logger.warning(f" ⚠️ [{self.platform}] 工具 {name} 失败: {e}") + return False, None, tools[-1][0] if tools else "none" +``` + +**关键语义**: +- `tool_level=1` → 只用 OpenCLI +- `tool_level=2` → OpenCLI → agent-browser(**默认**) +- `tool_level=3` → 三层全开 +- **第一个成功就停**;返回 `(成功?, 结果, 用了哪个工具)` + +--- + +## 四、适配器层(每平台一个) + +### 4.1 统一接口(`services/adapters/__init__.py` 的基类) + +**每个平台适配器必须实现 4 个方法**: + +```python +class ScrapeAdapter: + platform = "" # 子类覆写,如 "douyin" + adapter_name = "" + + async def collect_item(self, target, depth="light", tool_level=2) -> dict: + """抓单条内容详情""" + async def collect_user(self, user_id, limit=20) -> list[dict]: + """抓某个用户/作者的全部作品""" + async def collect_comments(self, item_id, limit=20) -> list[dict]: + """抓评论""" + async def collect_search(self, keyword, limit=20) -> list[dict]: + """按关键词搜索""" +``` + +**基类提供的公共能力**: +- `_try_tools(tool_level, funcs)` — 降级执行(见 3.4) +- `_run_opencli(args, timeout)` — 调 OpenCLI + 解析输出 +- `_parse_output(text)` — **输出格式三层兜底解析**(见下) +- `_parse_lines(text)` — 纯文本兜底解析 + +### 4.2 输出格式三层兜底(细节,容易踩坑) + +OpenCLI 的输出格式**不保证稳定**,所以解析做了三层: + +```python +def _parse_output(self, text: str): + # 1. JSON 检测(以 [ 或 { 开头)→ json.loads + if text.startswith("[") or text.startswith("{"): + try: + return json.loads(text) + except json.JSONDecodeError: + logger.warning("JSON 解析失败,尝试 YAML 兜底") + + # 2. YAML 解析(PyYAML 可用时) + try: + import yaml + parsed = yaml.safe_load(text) + if parsed is not None: + return parsed + except ImportError: + logger.debug("PyYAML 未安装,跳过 YAML") + + # 3. 纯文本兜底:按行解析 key:value + return self._parse_lines(text) +``` + +`_parse_lines` 甚至**专门处理了 `top_comments` 块**(评论在纯文本里的多行结构)。 + +**⛔ 移植注意**:如果目标环境没有 PyYAML,会静默降级到第三层 +(能跑,但嵌套结构会丢)。 + +### 4.3 各平台降级链(实测,**差异很大**) + +```python +# 抖音(douyin_scrape.py)—— 两条路径 +await self._try_tools(tool_level, [ + ("opencli", lambda: self._opencli_user_videos(...)), + ("agent-browser", lambda: self._browser_user_profile(...)), # ← 有浏览器降级 +]) + +# 小红书 / B站 / 知乎 —— ⚠️ 只有一条路径 +await self._try_tools(2, [ + ("opencli", lambda: self._opencli_user_notes(...)), # ← 没有降级! +]) + +# 通用网页(web_scrape.py)—— 不用 OpenCLI +await self._try_tools(tool_level, [ + ("web_crawler", lambda: self._web_crawl(target)), + ("agent-browser", lambda: self._browser_extract(target)), +]) +``` + +**对照表**: + +| 平台 | 降级链 | 有浏览器降级? | 备注 | +|---|---|---|---| +| **抖音** | `opencli` → `agent-browser` | ✅ | 唯一做了完整降级的平台 | +| **小红书** | `opencli`(单条) | ❌ | OpenCLI 挂了就抓不了 | +| **B站** | `opencli`(单条) | ❌ | 同上 | +| **知乎** | `opencli`(单条) | ❌ | 同上 | +| **通用网页** | `web_crawler` → `agent-browser` | ✅ | 走另一套(不用 OpenCLI) | + +**⚠️ 这是模块的真实局限**:三个平台**只有一条路径**,没有降级能力。 +接手时如果要提升健壮性,**最值得做的就是给它们补上浏览器降级** +(照抖音的 `_browser_*` 方法写即可)。 + +**另一个细节**:`collect_user` 传的是**硬编码的 `2`**(不是 `tool_level` 参数), +而 `collect_item` 才用传入的 `tool_level`: +```python +async def collect_user(self, user_id, limit=20): # 没有 tool_level 参数 + await self._try_tools(2, [...]) # ← 写死 2 +async def collect_item(self, target, depth, tool_level=2): + await self._try_tools(tool_level, [...]) # ← 用参数 +``` +→ **用户在前端设 `tool_level=1` 时,"抓用户主页"这条路仍会走到浏览器**。 + +### 4.35 登录态机制(★ 这是"如何用真实浏览器"的另一半) + +**问题**:抖音的数据接口需要登录态(cookie),怎么拿到? + +**答案**:`mediacrawler_adapter.py` 的做法 —— **用户在真实 Chrome 登录,程序通过 CDP 读解密后的 cookie**。 + +#### 为什么不能直接读 cookie 文件 + +代码注释原文(第 29 行): +``` +# Chrome 新版把 cookie 值加密存在 SQLite 里,必须通过 CDP 读解密后的值。 +``` + +Chrome v80+ 把 cookie **加密**存在 SQLite(`Cookies` 文件), +直接读文件拿到的是**密文**。必须让 **Chrome 自己解密** → 通过 **CDP 协议**问它。 + +#### 三步实现 + +**① 通过 CDP 读 cookie**(`_get_cookies`,第 68 行) +```python +async def _get_cookies() -> dict: + """从 CDP 连接读取 Chrome cookie(解密后的值)""" + ctx = await _ensure_cdp() # 确保 CDP 连接 + all_cookies = await ctx.cookies() # ← 让 Chrome 解密并返回 + ... +``` + +**② 转成 HTTP Header**(`_cookie_str`,第 86 行) +```python +def _cookie_str(cookies: dict) -> str: + return "; ".join(f"{k}={v}" for k, v in cookies.items()) +``` + +**③ 带 cookie 直调抖音 API**(`_http_get`,第 93 行) +```python +def _http_get(url: str, cookies: dict, timeout: int = 15) -> dict: + headers = { + "User-Agent": "...Chrome/150.0.0.0 Safari/537.36", + "Cookie": _cookie_str(cookies), # ★ + "Referer": "https://www.douyin.com/", + "Origin": "https://www.douyin.com", + } + ... +``` + +#### 登录态判断 + +```python +cookies = await _get_cookies() +has_session = bool(cookies.get("sessionid")) # ← 有 sessionid 就算已登录 +``` +对应 HTTP 端点 `GET /api/scrape/check-login`。 + +#### 让用户登录(巧妙的做法) + +代码注释(第 433 行): +```python +# ── 打开登录页(用 AppleScript 控制 Chrome,不需要 CDP) ── +"""在 Chrome 中打开抖音首页,让用户登录 + 登录后 cookie 自动保存到 Chrome profile,两个 Chrome 都会检测。""" +``` + +**→ 用 AppleScript 打开真实 Chrome**(不是 Playwright 控制的), +用户**在真实浏览器里手动登录** → cookie 存进 Chrome profile → +之后程序通过 CDP 读。 + +**这样最自然**:用户看到的是他熟悉的 Chrome,扫码登录, +✅ 不会被平台识别为"自动化登录"。 + +**⛔ 移植注意**:`open_login_page()` 用的是 **AppleScript**(`osascript`), +**macOS 专有**。Linux/Windows 要改实现(可以用 `open` 命令或直接 Playwright 打开)。 + +#### 一个隐藏技巧(避免暴露自动化) + +代码注释(第 227-229 行): +```python +# 获取热评(复用已有 Chrome 页面,不创建新标签页) +# Chrome 已有 douyin.com 页面,直接用它的 JS 上下文执行 fetch +# ⚠️ 不要 new_page() — 那会在 Chrome 中闪出新标签页 +``` + +**→ 复用用户已经打开的页面**执行 fetch,而不是新开标签页 +(新标签页会闪一下,且更易被识别)。 + +--- + +### 4.4 抖音的特有实现(两份代码,别搞混) + +**⚠️ 项目里有**两个**抖音相关模块**,职责不同: + +**① `services/adapters/douyin_scrape.py`(181 行)** +- 走**工具降级**(OpenCLI → 浏览器) +- 浏览器路径用 JS 从页面 **DOM 文本**提取数据: +```javascript +// _browser_video_page 的 JS(正则从页面文本抓「获赞/粉丝/关注」) +const body = document.body.innerText || ''; +const uidM = body.match(/抖音号[::]\s*(\S+)/); +const nickM = body.match(/@(\S+)/); +function extractNum(label) { + var m = body.match(new RegExp('(\\d+(?:\\.\\d+)?[万w]?)\\s*' + label)); + ... +} +return { aweme_id, title, author_nickname, douyin_id, digg_count, fans, following }; +``` + +**② `services/mediacrawler_adapter.py`(457 行)** +- **不走工具降级**,完全独立的实现 +- 文件头注释(原文): + > 全新架构:**不再依赖 CDP/Playwright/浏览器页面**。 + > 直接从 Chrome profile 读取 cookie,通过 HTTP 请求调用抖音 API。 + > 全程无窗口、无标签页、无闪烁。 +- 用于**追踪视频/作者**这类需要高频刷新的场景(详见子代理报告) + +**⛔ 关键澄清**:`mediacrawler_adapter.py` **虽然叫 mediacrawler,但不依赖 +MediaCrawler 这个开源项目**——它只是借用名字,实际是"读 Chrome cookie + HTTP 调 API"。 + +--- + +## 五、引擎层(`services/scrape_engine.py`,361 行) + +### 5.1 URL 解析(`resolve_target`,第 56 行) + +**7 类目标自动识别**(实测代码): + +```python +抖音短链 v.douyin.com/xxx → type=shortlink +抖音视频 douyin.com/video/{id} → type=video +抖音用户 douyin.com/user/{sec_uid} → type=user +小红书 xiaohongshu.com/explore/{id} → type=note +B站视频 bilibili.com/video/{BV} → type=video +B站用户 bilibili.com/space/{mid} → type=user +知乎 zhihu.com/answer/{id} | /question/ → type=item +通用网页 http(s)://... → type=page +纯 sec_uid MS4w 开头 或 len>20 → douyin/user +纯数字 len>=15 → douyin/video(aweme_id) +纯数字 其他 → zhihu/item +``` + +**短链解析**(`_resolve_shortlink`,第 136 行):用 `curl -sI` 拿 `Location` 头, +再用正则从跳转 URL 里抠出 `aweme_id`。 + +### 5.2 执行主流程(`run`,第 163 行) + +``` +run(request) + ├─ 1. resolve_urls(targets) → 标准化目标列表 + ├─ 2. 短链逐个解析 + ├─ 3. 判断同步 / 异步 + │ async_mode = request.async_mode 或 len(targets) > 50 + │ + ├─ 【异步分支】 + │ · run_id = uuid[:8] + │ · 内存状态 {status, total, completed, results, errors} + │ · asyncio.create_task(_run_async(...)) + │ · 立即返回 {status:"async", run_id} ← 前端轮询 + │ + └─ 【同步分支】 + · db.create_task("single", ...) + · for target: _scrape_one() → _save_item() + · db.update_task_status("completed", summary) + · 返回 {status, task_id, duration, total, success, errors, data} +``` + +### 5.3 单目标抓取(`_scrape_one`,第 261 行) + +```python +adapter = self._get_adapter(platform) # 按平台取适配器(带缓存) +if target["type"] == "user": + return await adapter.collect_user(target["target_id"]) # 返回 list +elif target["type"] in ("video", "note"): + return await adapter.collect_item(target["target_id"], depth, tool_level) +else: + return None +``` + +**返回值语义**(重要): +- `dict` → 单条内容 +- `list` → 多条(用户主页的所有作品) +- `None` → 失败 + +### 5.4 落库(`_save_item`,第 292 行) + +```python +db_id = self.db.insert_item( + task_id=..., platform=..., item_id=..., url=..., title=..., + author_name=..., author_id=..., published_at=..., text_content=..., + tags=..., stats=..., extra=..., media=...) + +comments = item.get("comments", []) +if comments and db_id: + self.db.insert_comments(db_id, comments) # ★ 评论独立表 +``` + +### 5.5 异步模式(`_run_async`,第 315 行) + +- ✅ **同时写内存 + 落库**(内存态供轮询,落库供持久化) +- 内存态在 `self._tasks[run_id]`(**进程重启即丢**) +- 查询用 `get_async_result(run_id)` + +### 5.6 ⚠️ 已知的局限(接手时要清楚) + +| 局限 | 说明 | +|---|---| +| **没有限流/并发控制** | `for target in ready:` 是**纯串行**,目标多时会慢;也没有请求间隔(可能触发平台风控) | +| **异步态存内存** | 重启 Dashboard 后 `_tasks` 丢失(但库里有记录,前端看不到进度) | +| **无重试** | 单目标失败只记 `errors`,不重试 | +| **adapter 实例缓存** | `self._adapters` 进程内缓存(无清理) | + +--- + +### 5.7 一个完整请求的生命周期(跟着走一遍最快懂) + +以"采集某抖音作者的全部视频"为例: + +``` +① 前端 + POST /api/scrape/run + {"targets": ["r606391422378804368"], "tool_level": 2} + │ + ▼ +② routes/scrape.py:44 api_scrape_run() + 组装 request → engine.run(request) + │ + ▼ +③ scrape_engine.py:163 run() + ├─ resolve_urls(["r6063..."]) + │ → resolve_target() 识别:以 MS4w 开头 → 抖音 sec_uid + │ → [{"platform":"douyin","type":"user","target_id":"r6063...","status":"resolved"}] + │ + ├─ 同步或异步?len(targets)=1,不大于 50 → 同步 + │ + ├─ db.create_task("single","douyin",...) → task_id = 1 + │ + ├─ for target: _scrape_one(target,"douyin","light",2) + │ │ + │ ▼ + │ scrape_engine.py:261 + │ _get_adapter("douyin") → DouyinScrapeAdapter() (进程内缓存) + │ type=="user" → adapter.collect_user("r6063...") + │ │ + │ ▼ + │ douyin_scrape.py:21 collect_user() + │ _try_tools(2, [("opencli", ...), ("agent-browser", ...)]) + │ │ + │ ├─ 尝试 1:_opencli_user_videos() + │ │ _run_opencli(["douyin","user-videos","r6063...", + │ │ "--limit","20","--with_comments","true","-f","json"]) + │ │ → subprocess 执行 opencli(Node CLI) + │ │ → _parse_output() ← JSON → YAML → 纯文本 三层兜底 + │ │ → 成功返回 list[dict] → _try_tools 立刻返回,不再降级 + │ │ + │ └─ 尝试 1 失败(OpenCLI 没装/超时/报错) + │ → 尝试 2:_browser_user_profile() + │ → Playwright 启真实 Chrome → 打开页面 → JS 提取 + │ + │ → 每条数据 _to_schema() 转成统一格式 + │ → 返回 list + │ + ├─ for item: _save_item(task_id=1, item) + │ db.insert_item(...) → db_id((platform,item_id) 唯一,重复则忽略) + │ db.insert_comments(db_id, comments) ← 评论另存 + │ + ├─ db.update_task_status(1, "completed", summary={success:N, errors:0}) + │ + └─ return {status:"completed", task_id:1, duration, total, success, errors, data:[...]} + │ + ▼ +④ 前端拿到 data,渲染列表 +``` + +**异步分支的差异**(目标 > 50 个,或显式 `async_mode=true`): +``` +run() 立即返回 {status:"async", run_id:"ab12cd34"} + ↓(后台) +asyncio.create_task(_run_async(run_id, targets, ...)) + ↓ + 建持久化任务 → 逐个 _scrape_one + _save_item + ↓ + 进度写内存 self._tasks[run_id](供轮询) + 结果写 SQLite(供持久化) + ↓ +前端轮询 get_async_result(run_id) 看进度 +``` + +**⚠️ 注意**:异步进度**只在内存**,Dashboard 重启就丢 +(库里数据还在,但前端看不到进度了)。 + +--- + +## 六、存储层(`services/scrape_db.py`,415 行) + +### 6.1 数据库位置 + +```python +DEFAULT_DB = AGENT_LOCAL / "data" / "scrape.db" +``` +(`AGENT_LOCAL` 是环境变量;默认 `~/workbuddy-agent-os/agent-local`) + +### 6.2 四张表(实测 `CREATE TABLE`) + +```sql +-- ① 采集任务(一次 run 一条) +scrape_tasks( + id, type, -- single / batch / scheduled + platform, target, -- 目标(批量时是 JSON 数组) + depth, tool_level, machine, + status, -- pending / running / completed / failed + total_targets, completed_targets, + summary, -- 摘要 JSON + created_at +) + +-- ② 采集到的内容 +scrape_items( + id, task_id → scrape_tasks, + platform, item_id, -- 平台内唯一 ID + url, title, author_name, author_id, + published_at, collected_at, text_content, tags, + ... -- 还有 stats / extra / media 等 +) + +-- ③ 评论 +scrape_comments(id, item_db_id → scrape_items, + author_name, text, likes, replied_at) + +-- ④ 采集源(长期跟踪) +scrape_sources( + id, platform, source_type, -- user / hashtag / keyword / url_list / api + target, display_name, + category, -- 自定义分类 + notes, + schedule, -- CRON(定期采集) + depth, tool_level, last_collected, + status -- active / paused +) +``` + +### 6.3 方法清单(22 个,实测) + +``` +任务:create_task / update_task_status / get_task / list_tasks +内容:insert_item / get_item_id / item_exists / get_item / list_items +评论:insert_comments / get_comments +采集源:upsert_source / update_source / list_sources / get_due_sources / + update_source_collected / delete_source +统计:count_by_platform / count_today / sources_count / task_stats +``` + +### 6.4 去重机制 + +`insert_item` **依赖 `(platform, item_id)` 唯一约束** —— 重复插入时 +用 `INSERT OR IGNORE` 模式(验证文档 L1-2 有测例)。 + +--- + +## 七、HTTP 层(`routes/scrape.py`,711 行 / 22 端点) + +### 7.1 端点清单(实测) + +``` +采集 + POST /api/scrape/run 发起采集(targets + depth + tool_level) + POST /api/scrape/resolve 只解析 URL,不采集 + POST /api/scrape/douyin-stats 抖音数据查询 + +查询 + GET /api/scrape/title 取标题 + GET /api/scrape/result 结果 + GET /api/scrape/tasks 任务列表 + GET /api/scrape/items 内容列表 + GET /api/scrape/items/{id} 单项详情 + GET /api/scrape/stats 统计 + +采集源管理 + POST /api/scrape/sources 新建源 + GET /api/scrape/sources 源列表 + DEL /api/scrape/sources/{id} 删源 + +追踪(视频 / 作者) + POST /api/scrape/track-video 追踪视频 + GET /api/scrape/tracked-videos 已追踪视频 + POST /api/scrape/delete-tracked/{id} + POST /api/scrape/refresh-video/{id} 刷新单个视频 + POST /api/scrape/track-author 追踪作者 + GET /api/scrape/tracked-authors 已追踪作者 + POST /api/scrape/refresh-author/{id} + GET /api/scrape/author-history/{id} + +主题 / 登录 + POST /api/scrape/import-topics 批量导入主题 + GET /api/scrape/check-login 检测登录态 + GET /api/scrape/open-login 打开登录 +``` + +### 7.2 关键实现(实测) + +**`POST /api/scrape/run`**(第 43 行)—— 前端发起采集的唯一入口: +```python +@router.post("/run") +async def api_scrape_run(data: dict = {}): + targets = data.get("targets", data.get("target", [])) + if isinstance(targets, str): + targets = [targets] # 兼容单个字符串 + request = { + "targets": targets, + "platform": data.get("platform", "auto"), + "depth": data.get("depth", "light"), + "tool_level": data.get("tool_level", 2), # ← 默认 2(OpenCLI + 浏览器) + "machine": data.get("machine", ""), + "multi_machine": data.get("multi_machine", False), + "async_mode": data.get("async_mode", False), + } + engine = _get_engine() # 模块级单例 + result = await engine.run(request) + return {"status": "ok", **result} +``` + +**登录态两端点**(第 692 / 703 行)—— 都委托给 `mediacrawler_adapter`: +```python +GET /api/scrape/check-login → mediacrawler_adapter.check_login_status() +POST /api/scrape/open-login → mediacrawler_adapter.open_login_page() +``` + +### 7.3 引擎实例 + +路由层用**模块级单例**拿 engine(`_get_engine()`), +所以 `engine._tasks`(异步态)和 `engine._adapters`(适配器缓存) +在整个 Dashboard 进程内共享。 + +--- + +## 七·五、Dashboard 插件(`plugins/crawl.py`,91 行) + +抓取模块**作为 Dashboard 插件**注册(提供概览统计,不是核心逻辑): + +```python +class CrawlDashboardPlugin(DashboardPlugin): + name = "crawl" + label = "内容抓取" + icon = "📡" + order = 35 +``` + +**它做三件事**: +1. `summary()` — 概览:总抓取数 / 今日新增 / 抓取源(从 `ScrapeDB` 读) +2. `detail(machine)` — 指定机器的详情 +3. `actions()` — 快捷操作(跳转 `scrape` 视图) + +**注意**:它有个**兜底设计** —— 如果 `ScrapeDB` 不可用(数据库损坏/权限), +会退化成**统计知识库里的 md 文件数**,而不是报错。 + +**⛔ 移植注意**:如果目标项目没有这套插件框架,`plugins/crawl.py` +可以直接丢弃(它只是 Dashboard 的展示层,不影响抓取功能本身)。 + +--- + +## 八、前端(`frontend/src/views/scrape.js`,131 KB) + +⚠️ **这是模块里最大的单文件**(131 KB)。功能覆盖:22 个端点的界面。 + +**建议接手团队**:不要照搬这个前端,按第七节的 HTTP 契约重写。 +理由:131 KB 单文件难维护,且和本项目的视图框架耦合。 + +--- + +## 九、验证方案(项目里已有现成的) + +`services/scrape_validation.md`(314 行)已经写了 **7 个 Level 的验证清单**: + +``` +L0 基础设施(3 项) Python import / SQLite 建库 / FastAPI 路由注册 +L1 数据库层(3 项) 建任务 / 写入+去重 / 评论入库 +L2 适配器 Mock(3 项) 工具降级逻辑 / 一级失败二级成功 +L3 适配器真实(3 项) 抖音用户采集 / 小红书 / 详情+评论 ← 需 OpenCLI + 登录态 +L4 引擎层(3 项) 解析 URL / 执行采集 / 异步轮询 +L5 API 层(4 项) curl 打 4 个端点 +L6 前端(4 项) 浏览器里操作 +L7 异常(1+ 项) OpenCLI 不可用时应抛清晰错误 +``` + +**⚠️ 但要注意**:该文档写于 2026-07-16,**里面的方法名已过时**: +``` +文档写 resolve_targets ← 不存在 +代码里是 resolve_urls ← 实际 +文档写 get_result ← 不存在 +代码里是 get_async_result ← 实际 +``` +**以代码为准**。 + +--- + +## 十、移植清单(换环境要改什么) + +| # | 依赖 | 位置 | 处理 | +|---|---|---|---| +| 1 | **OpenCLI**(Node CLI) | `adapters/__init__.py` 的 `OPENCLI` 常量 | 目标机 `npm i -g @jackwener/opencli`,或设 `OPENCLI_PATH` | +| 2 | **真实 Chrome** | `browser_helpers.py` 的 `_CHROME_PATHS` | Linux/Windows 要改路径(如 `/usr/bin/google-chrome`) | +| 3 | **Playwright** | pip | `pip install playwright && playwright install chromium` | +| 4 | **PyYAML** | pip(可选但强烈建议) | 不装会降级到纯文本解析(丢嵌套结构) | +| 5 | **AGENT_LOCAL** 环境变量 | `scrape_db.py:18` | 定义了才能定位 `scrape.db` | +| 6 | **平台登录态** | Chrome profile | 目标机需手动登录一次目标平台 | +| 7 | `SCRAPE_CHROME_USER_DATA` | 环境变量(可选) | 要持久化登录态时指定 | + +### 最小可运行子集(只要"能采集") + +``` +services/scrape_db.py 数据库(4 表) +services/adapters/__init__.py 基类 + 降级 + OpenCLI 调用 + 输出解析 +services/adapters/browser_helpers.py 浏览器降级 +services/adapters/<目标平台>_scrape.py 目标平台适配器 +``` +—— 这 4 个文件就能跑通单平台采集,不需要 engine/routes/前端。 + +--- + +## 十一、接手建议(按顺序) + +``` +第 1 步 装 OpenCLI + Playwright + 真实 Chrome,跑 validation.md 的 L0/L1 + (这两级零外部依赖,能验证环境对不对) + +第 2 步 跑 L2(Mock 测试)—— 验证降级逻辑,不需要真实平台 + 这时你已经能理解 _try_tools 的语义 + +第 3 步 登录目标平台,跑 L3(真实采集)—— 第一次真正拿数据 + 如果 OpenCLI 不通,会看到它降级到浏览器,日志里有 ⚠️ + +第 4 步 跑 L4/L5(引擎 + API) + +第 5 步 替换前端(不要照搬 131 KB) + +第 6 步 加你要的东西:限流 / 重试 / 并发控制(现在都没有) +``` + +--- + +## 十二、这个模块还没做的事(接手可以补) + +``` +① 限流与请求间隔 —— 现在纯串行、无间隔,目标多时可能触发平台风控 +② 失败重试 —— 现在失败只记 errors +③ 并发控制 —— 没有信号量,大量目标只能串行 +④ 异步态持久化 —— run 进度存内存,重启即丢 +⑤ 登录态自动检测 —— check-login 端点有,但采集前没强制校验 +⑥ 代理支持 —— 没有看到代理配置(多账号场景会需要) +``` + +--- + +## 附录:本说明的取证方式(可复现) + +```bash +cd 05_tools/10_dashboard +# 架构 +head -20 services/adapters/__init__.py +# 浏览器("套真实浏览器"的做法) +sed -n '1,70p' services/adapters/browser_helpers.py +# 降级机制 +grep -n "_try_tools" -A 14 services/adapters/__init__.py +# 引擎流程 +grep -nE "^ (async )?def " services/scrape_engine.py +# 表结构 +grep -n "CREATE TABLE" -A 12 services/scrape_db.py +# 端点 +grep -nE "^@router\." routes/scrape.py +``` diff --git a/api/monitor/douyin_api.py b/api/monitor/douyin_api.py index 851ae40..9fc6f2f 100644 --- a/api/monitor/douyin_api.py +++ b/api/monitor/douyin_api.py @@ -69,6 +69,12 @@ COMMENT_PATH = "/aweme/v1/web/comment/list/" # 一次请求的超时。抖音这两个接口正常都在一秒内返回。 REQUEST_TIMEOUT_SECONDS = 20.0 + +# 问浏览器要 UA / client hints 的超时。**这个必须有。** +# ``page.evaluate`` 打在一个渲染进程已经卡住的标签页上会**永远不返回**,而问身份是采集的 +# 第一步 —— 它一挂,整个 run 就永远停在「运行中」(真踩过:标签页 URL 是空的, +# cookies() 正常,evaluate 一直不回来)。 +EVALUATE_TIMEOUT_SECONDS = 8.0 # 单页最多要多少条。接口自己有上限,要多了也没用。 MAX_PAGE_SIZE = 20 @@ -117,6 +123,61 @@ _IDENTITY_TTL_SECONDS = 120.0 _identity_cache: Optional[Tuple[float, BrowserIdentity]] = None +async def _safe_evaluate(page: Any, expression: str) -> Any: + """在页面上求值,带超时;任何失败都返回 None。 + + **不要直接调 ``page.evaluate``** —— 在渲染进程卡住的标签页上它会永远不返回(见 + ``EVALUATE_TIMEOUT_SECONDS`` 那段)。 + """ + try: + return await asyncio.wait_for( + page.evaluate(expression), timeout=EVALUATE_TIMEOUT_SECONDS + ) + except Exception: + return None + + +async def _identity_from_pages(context: Any) -> Tuple[str, Dict[str, str]]: + """问出 UA 和 client hints。 + + 不假设第一个标签页是好的 —— 它可能停在 URL 为空、渲染进程已卡住的状态(实测过)。 + 所以逐个试、每个都带超时;优先抖音页面,全都不行就临时开一个干净页问完关掉。 + + 拿不到就返回空 —— 调用方据此退回库里那份 cookie,而不是拿一组编出来的指纹去请求 + (那比没有更糟,见 BrowserIdentity 的说明)。 + """ + from media_platform.douyin.help import client_hint_headers + + pages = list(context.pages) + pages.sort(key=lambda page: 0 if "douyin" in (page.url or "") else 1) + for page in pages: + user_agent = await _safe_evaluate(page, "() => navigator.userAgent") + if user_agent: + hints = client_hint_headers( + await _safe_evaluate(page, "() => navigator.userAgentData || null") + ) + return user_agent, hints or {} + + temp = None + try: + temp = await asyncio.wait_for( + context.new_page(), timeout=EVALUATE_TIMEOUT_SECONDS + ) + user_agent = await _safe_evaluate(temp, "() => navigator.userAgent") + hints = client_hint_headers( + await _safe_evaluate(temp, "() => navigator.userAgentData || null") + ) + return user_agent or "", hints or {} + except Exception: + return "", {} + finally: + if temp is not None: + try: + await temp.close() + except Exception: + pass + + async def _read_browser() -> Optional[BrowserIdentity]: """连上 CDP 浏览器,一次取齐 cookie、UA、client hints。 @@ -134,14 +195,12 @@ async def _read_browser() -> Optional[BrowserIdentity]: return None # contexts[0] 是真实 profile。**不要 new_context()** —— 那是无痕式的,读不到登录态。 context = browser.contexts[0] - cookies = await context.cookies() + cookies = await asyncio.wait_for( + context.cookies(), timeout=EVALUATE_TIMEOUT_SECONDS + ) # UA 和 hints 要从页面里问 —— 它们是浏览器自己的事实,写死迟早对不上。 - page = context.pages[0] if context.pages else await context.new_page() - user_agent = await page.evaluate("() => navigator.userAgent") - hints = client_hint_headers( - await page.evaluate("() => navigator.userAgentData || null") - ) + user_agent, hints = await _identity_from_pages(context) except Exception as exc: utils.logger.warning(f"[douyin_api] 读浏览器身份失败:{exc}") return None diff --git a/api/monitor/runner.py b/api/monitor/runner.py index ad0906e..cd51dae 100644 --- a/api/monitor/runner.py +++ b/api/monitor/runner.py @@ -221,17 +221,31 @@ async def execute_task(task_id: int, trigger: str = "manual") -> IngestResult: in_process_tail: List[str] = [] if platform == adapters.PLATFORM_DY: - fetched = await douyin_fetch.collect( - out_dir, - platform=platform, - mode=mode, - limit=max_notes_count, - want_comments=enable_comments, - comment_limit=max_comments_count, - targets=targets, - known_aweme_ids=known_aweme_ids, - cookie=cookie, - ) + try: + # 也要有超时。爬虫那条路靠 run_and_wait(timeout=...) 兜底,这条路没有子进程、 + # 没人管 —— 里面**任何一次卡住都会让 run 永远停在「运行中」**(真踩过: + # page.evaluate 打在一个卡死的标签页上不返回)。 + fetched = await asyncio.wait_for( + douyin_fetch.collect( + out_dir, + platform=platform, + mode=mode, + limit=max_notes_count, + want_comments=enable_comments, + comment_limit=max_comments_count, + targets=targets, + known_aweme_ids=known_aweme_ids, + cookie=cookie, + ), + timeout=timeout_seconds, + ) + except asyncio.TimeoutError: + fetched = { + "notes": 0, + "comments": 0, + "errors": [f"抖音采集超过 {timeout_seconds} 秒仍未完成,已放弃这一轮"], + "jsonl_dir": "", + } in_process_tail = list(fetched["errors"]) # 一条都没采到 = 这一轮失败,并把**真因**当作退出诊断传下去。否则它会掉进 # ingest 的「疑似登录失效」分支 —— 又骗人一次,正是这套东西一直在犯的毛病。 diff --git a/tests/test_douyin_api.py b/tests/test_douyin_api.py index 3feaa50..aef3cc0 100644 --- a/tests/test_douyin_api.py +++ b/tests/test_douyin_api.py @@ -116,6 +116,101 @@ class TestNormalizeAweme: assert record["create_time"] == 0 +class TestIdentityFromPages: + """身份得从浏览器里问,但**不能被一个卡死的标签页拖住**。 + + 实测过:标签页 URL 为空、渲染进程卡死,``page.evaluate`` 永远不返回;而问身份是采集的 + 第一步 —— 没超时的话整轮就挂在那儿,run 永远停在「运行中」。 + """ + + class _Page: + def __init__(self, url, *, user_agent=None, hang=False): + self.url = url + self._user_agent = user_agent + self._hang = hang + self.closed = False + + async def evaluate(self, expression): + if self._hang: + await asyncio.sleep(30) # 模拟渲染进程卡死 + if expression.startswith("() => navigator.userAgentData"): + return { + "brands": [{"brand": "Chrome", "version": "155"}], + "mobile": False, + "platform": "Linux", + } + return self._user_agent + + async def close(self): + self.closed = True + + class _Context: + def __init__(self, pages, temp=None): + self.pages = pages + self._temp = temp + self.made_temp = False + + async def new_page(self): + self.made_temp = True + if self._temp is None: + raise AssertionError("这个用例不该走到临时页") + return self._temp + + def _fast_timeout(self, monkeypatch): + monkeypatch.setattr(douyin_api, "EVALUATE_TIMEOUT_SECONDS", 0.05) + + @pytest.mark.asyncio + async def test_a_hanging_page_is_skipped(self, monkeypatch): + self._fast_timeout(monkeypatch) + stuck = self._Page("", hang=True) + good = self._Page("https://example.com/", user_agent="UA-of-good-page") + + user_agent, hints = await douyin_api._identity_from_pages( + self._Context([stuck, good]) + ) + + assert user_agent == "UA-of-good-page" + assert hints["sec-ch-ua"] == '"Chrome";v="155"' + + @pytest.mark.asyncio + async def test_a_page_without_a_user_agent_is_skipped(self, monkeypatch): + self._fast_timeout(monkeypatch) + blank = self._Page("", user_agent=None) + good = self._Page("https://example.com/", user_agent="UA-of-good-page") + + user_agent, _ = await douyin_api._identity_from_pages( + self._Context([blank, good]) + ) + + assert user_agent == "UA-of-good-page" + + @pytest.mark.asyncio + async def test_it_opens_a_temporary_page_when_nothing_else_works(self, monkeypatch): + self._fast_timeout(monkeypatch) + stuck = self._Page("", hang=True) + temp = self._Page("about:blank", user_agent="UA-of-temp-page") + context = self._Context([stuck], temp=temp) + + user_agent, _ = await douyin_api._identity_from_pages(context) + + assert context.made_temp is True + assert user_agent == "UA-of-temp-page" + assert temp.closed is True, "临时页问完要关掉,别在操作者的浏览器里留垃圾" + + @pytest.mark.asyncio + async def test_a_douyin_page_is_preferred(self, monkeypatch): + """有抖音页面就先问它 —— 它才是我们要模仿的那个身份。""" + self._fast_timeout(monkeypatch) + other = self._Page("https://example.com/", user_agent="UA-of-other") + douyin = self._Page("https://www.douyin.com/explore", user_agent="UA-of-douyin") + + user_agent, _ = await douyin_api._identity_from_pages( + self._Context([other, douyin]) + ) + + assert user_agent == "UA-of-douyin" + + class TestGet: """`_get` 的失败路径 —— 它们决定了失败会不会被伪装成「这个博主没作品」。""" diff --git a/tests/test_monitor_runner.py b/tests/test_monitor_runner.py index 9ad933a..d853faa 100644 --- a/tests/test_monitor_runner.py +++ b/tests/test_monitor_runner.py @@ -5,6 +5,8 @@ 约定它一个都不沾。凡是写在那里面的东西,这条路都得单独有一份。 """ +import asyncio + import pytest import pytest_asyncio from sqlalchemy import select @@ -80,3 +82,29 @@ class TestDouyinRunStatus: await runner_module.execute_task(task_id, trigger="manual") assert seen["status"] == RUN_RUNNING + + @pytest.mark.asyncio + async def test_a_hanging_collect_does_not_leave_the_run_running(self, db, monkeypatch): + """进程内那条路也要有超时。 + + 爬虫那条靠 ``run_and_wait(timeout=...)`` 兜底,这条路没有子进程、没人管 —— + 里面任何一次卡住(实测过 ``page.evaluate`` 打在一个卡死的标签页上不返回)都会让 + run 永远停在「运行中」,界面上看起来就是任务卡死了。 + """ + task_id = await _make_douyin_task() + async with monitor_db.get_session() as session: + task = await session.get(MonitorTask, task_id) + task.run_timeout_seconds = 1 # 把超时压到 1 秒,别让测试真等 + + async def hanging_collect(out_dir, **kwargs): + await asyncio.sleep(60) + raise AssertionError("不该走到这里") + + monkeypatch.setattr(runner_module.douyin_fetch, "collect", hanging_collect) + + await runner_module.execute_task(task_id, trigger="manual") + + async with monitor_db.get_session() as session: + run = await session.scalar(select(MonitorRun).order_by(MonitorRun.id)) + assert run.status != RUN_RUNNING + assert "超时" in (run.error_message or "") or "超过" in (run.error_message or "")