fix(monitor): 抖音任务永远「运行中」—— page.evaluate 卡在一个死掉的标签页上,而我没给超时
用户报「一直在运行中」。库里那条 run 是真的卡住了:日志里一条采集输出都没有,
说明它卡在 collect 里、还没走到任何日志。探针定位到:
标签页: [''] ← 一个 URL 为空的标签页,渲染进程已卡死
context.cookies(): OK, 90 个 ← cookie 读得到
page.evaluate('navigator.userAgent'): **永远不返回**
而问浏览器要身份(UA + client hints)是采集的**第一步**,`page.evaluate` 又**没设超时** ——
于是整轮挂在那儿,run 永远停在「运行中」。
三处修复,各挡一层:
1. `page.evaluate` / `context.cookies()` 全部加超时(8 秒)。卡住就跳过,不再无限等。
2. 不假设第一个标签页是好的:逐个试、优先抖音页;全都不行就临时开一个干净页问完关掉。
拿不到就退回库里那份 cookie —— **不编造指纹**,那比没有更糟。
3. **进程内那条路补上整体超时**:爬虫那条靠 `run_and_wait(timeout=...)` 兜底,这条路
没有子进程、没人管,里面任何一次卡住都会变成永久的「运行中」。
测试 +6:卡死的页会被跳过(真 sleep,验的正是超时)、没 UA 的页跳过、全不行时开临时页
并关掉它、优先抖音页;以及整轮卡住时 run 不会停在 running(含超时原因)。
This commit is contained in:
@@ -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 <sec_uid> --limit 20 --with_comments true -f json
|
||||
opencli douyin stats <aweme_id> -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
|
||||
```
|
||||
@@ -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
|
||||
|
||||
+25
-11
@@ -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 的「疑似登录失效」分支 —— 又骗人一次,正是这套东西一直在犯的毛病。
|
||||
|
||||
@@ -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` 的失败路径 —— 它们决定了失败会不会被伪装成「这个博主没作品」。"""
|
||||
|
||||
|
||||
@@ -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 "")
|
||||
|
||||
Reference in New Issue
Block a user