Files
MediaCrawler/tests/test_monitor_ingest.py
butubb 4e60524f37
Deploy VitePress site to Pages / build (push) Canceled after 0s
Deploy VitePress site to Pages / Deploy (push) Canceled after 0s
feat: 监控面板 / 登录鉴权 / 多平台切换 / MySQL
在上游 MediaCrawler 之上新增一层:

- 监控层 api/monitor/ —— 多博主/多笔记的定时采集、指标快照差分、报表、
  企业微信通知。每轮采集写入独立目录,差分才成立。
- WebUI 登录鉴权 api/auth.py —— PBKDF2 口令 + 服务端会话,/api 全接口防护。
  WebSocket 单独加依赖:BaseHTTPMiddleware 对 ws 作用域直接放行,覆盖不到。
- 全局平台切换 + 能力矩阵 —— 如实区分「爬虫模块支持」与「监控层已接线」,
  未接通的平台直接拒绝建任务,而不是静默跑空。
- 监控库改用 MySQL 5.7(可回退 SQLite 供测试):逐表强制 utf8mb4
  (服务端与库默认都是 latin1),启动校验所连 schema 以防写错库,
  连接池 recycle + pre_ping 应对 MySQL 的 8 小时空闲断连。

修复上游缺陷:

- xhs/core.py: 主页抓取失败会跳掉整个博主,导致一条作品都抓不到,
  而那份资料只喂给一个空函数。改为尽力而为,失败不中断。
- xhs/login.py: cookie 登录只注入 web_session,冷启动签名会失败。
  新增 INJECT_ALL_COOKIES 开关(默认关闭,原有行为不变)。
- requirements.txt: 补上 websockets。它在上游 pyproject.toml 里有声明、
  这里漏了,导致 uvicorn 没有 WebSocket 能力,实时日志流从未工作。

改动过的上游文件清单及合并方式见 UPSTREAM.md。

测试:492 passed(另有 1 个既有的 Windows/gbk 上游测试失败,与本改动无关)
2026-10-07 09:58:40 +08:00

528 lines
19 KiB
Python

# -*- coding: utf-8 -*-
# Copyright (c) 2025 [email protected]
#
# This file is part of MediaCrawler project.
# Repository: https://github.com/NanmiCoder/MediaCrawler/blob/main/tests/test_monitor_ingest.py
# GitHub: https://github.com/NanmiCoder
# Licensed under NON-COMMERCIAL LEARNING LICENSE 1.1
#
# 声明:本代码仅供学习和研究目的使用。使用者应遵守以下原则:
# 1. 不得用于任何商业用途。
# 2. 使用时应遵守目标平台的使用条款和robots.txt规则。
# 3. 不得进行大规模爬取或对平台造成运营干扰。
# 4. 应合理控制请求频率,避免给目标平台带来不必要的负担。
# 5. 不得用于任何非法或不当的用途。
#
# 详细许可条款请参阅项目根目录下的LICENSE文件。
# 使用本代码即表示您同意遵守上述原则和LICENSE中的所有条款。
"""Offline tests for the monitoring ingest/diff layer.
These run without network, browser or login and cover the correctness caveats
that matter most: baseline suppression, count parsing, NULL-vs-zero, the
posted/seen comment split, idempotency, and the silent-cookie-failure signal.
"""
import json
from pathlib import Path
from typing import Any, Dict, List, Optional
import pytest
import pytest_asyncio
from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker, create_async_engine
from sqlalchemy.pool import StaticPool
from tools.time_util import get_current_timestamp
from api.monitor.ingest import describe_exit_code, ingest_run, parse_count
from api.monitor.models import (
EVENT_AUTH_FAILURE,
EVENT_METRIC_DELTA,
EVENT_NEW_COMMENT_POSTED,
EVENT_NEW_COMMENT_SEEN,
EVENT_NEW_NOTE,
EVENT_NO_DATA,
EVENT_RUN_FAILED,
MODE_CREATOR,
MonitorBase,
MonitorEvent,
MonitorNote,
MonitorNoteMetric,
MonitorRun,
MonitorTask,
RUN_FAILED,
RUN_PARTIAL,
RUN_SUCCESS,
)
@pytest_asyncio.fixture
async def db():
"""An isolated in-memory monitoring database."""
engine = create_async_engine("sqlite+aiosqlite://", poolclass=StaticPool)
async with engine.begin() as conn:
await conn.run_sync(MonitorBase.metadata.create_all)
factory = async_sessionmaker(engine, class_=AsyncSession, expire_on_commit=False)
async with factory() as db_session:
yield db_session
await engine.dispose()
async def _make_task(db: AsyncSession, **overrides) -> MonitorTask:
defaults = dict(
name="test task",
platform="xhs",
mode=MODE_CREATOR,
enabled=True,
interval_minutes=60,
max_notes_count=20,
enable_comments=True,
max_comments_count=50,
run_timeout_seconds=3600,
created_at=0,
updated_at=0,
)
defaults.update(overrides)
task = MonitorTask(**defaults)
db.add(task)
await db.flush()
return task
async def _make_run(
db: AsyncSession,
task: MonitorTask,
started_at: int,
exit_code: Optional[int] = 0,
) -> MonitorRun:
run = MonitorRun(
task_id=task.id,
trigger="manual",
status=RUN_SUCCESS,
phase=task.mode,
save_data_path="",
queued_at=started_at,
not_before=0,
started_at=started_at,
exit_code=exit_code,
)
db.add(run)
await db.flush()
return run
def _write_run_dir(
root: Path,
notes: List[Dict[str, Any]],
comments: Optional[List[Dict[str, Any]]] = None,
) -> Path:
"""Write a run's jsonl output in the crawler's own layout."""
jsonl_dir = root / "xhs" / "jsonl"
jsonl_dir.mkdir(parents=True, exist_ok=True)
contents = jsonl_dir / "creator_contents_2026-01-01.jsonl"
contents.write_text(
"\n".join(json.dumps(n, ensure_ascii=False) for n in notes),
encoding="utf-8",
)
if comments is not None:
comment_file = jsonl_dir / "creator_comments_2026-01-01.jsonl"
comment_file.write_text(
"\n".join(json.dumps(c, ensure_ascii=False) for c in comments),
encoding="utf-8",
)
return root
def _note(note_id: str, liked: Any = "10", **extra) -> Dict[str, Any]:
record = {
"note_id": note_id,
"title": f"title-{note_id}",
"note_url": f"https://www.xiaohongshu.com/explore/{note_id}",
"image_list": "https://img/cover.jpg",
"creator_hash": "hash",
"time": 1700000000000,
"liked_count": liked,
"comment_count": "1",
"collected_count": "1",
"share_count": "1",
}
record.update(extra)
return record
def _comment(comment_id: str, note_id: str, create_time: int, **extra) -> Dict[str, Any]:
record = {
"comment_id": comment_id,
"note_id": note_id,
"content": f"content-{comment_id}",
"nickname": "u***r",
"creator_hash": "hash",
"create_time": create_time,
"like_count": "0",
"sub_comment_count": 0,
"parent_comment_id": "",
}
record.update(extra)
return record
async def _events(db: AsyncSession, event_type: Optional[str] = None) -> List[MonitorEvent]:
stmt = select(MonitorEvent)
if event_type:
stmt = stmt.where(MonitorEvent.type == event_type)
return list((await db.scalars(stmt)).all())
# --------------------------------------------------------------------------
# parse_count
# --------------------------------------------------------------------------
class TestParseCount:
@pytest.mark.parametrize(
"raw,expected",
[
("1234", 1234),
("1.2万", 12000),
("1.2w", 12000),
("3亿", 300000000),
("1,234", 1234),
(42, 42),
],
)
def test_parses_platform_formats(self, raw, expected):
assert parse_count(raw) == expected
@pytest.mark.parametrize("raw", ["", None, "暂无", "-", "abc", True])
def test_unparseable_values_return_none(self, raw):
assert parse_count(raw) is None
# --------------------------------------------------------------------------
# Exit codes
# --------------------------------------------------------------------------
class TestExitCodeStorage:
"""Guards a bug that only showed up when the data moved to MySQL.
Windows reports process failures as unsigned 32-bit NTSTATUS values
(0xC0000142 = 3221225794). That overflows MySQL's signed INT, while SQLite's
dynamic typing accepted it happily -- so the column silently worked until a
real migration hit it with real data.
"""
def test_column_is_bigint_not_int(self):
from sqlalchemy import BigInteger
from api.monitor.models import MonitorRun
column_type = MonitorRun.__table__.c.exit_code.type
assert isinstance(column_type, BigInteger), (
f"exit_code must be BigInteger to hold unsigned 32-bit codes, got {column_type!r}"
)
@pytest.mark.asyncio
async def test_an_ntstatus_value_round_trips(self, db):
task = await _make_task(db)
run = await _make_run(db, task, started_at=1000, exit_code=3221225794)
await db.commit()
stored = await db.scalar(
select(MonitorRun.exit_code).where(MonitorRun.id == run.id)
)
assert stored == 3221225794
class TestDescribeExitCode:
def test_windows_status_code_is_decoded(self):
"""3221225794 is 0xC0000142, which is meaningless without decoding."""
message = describe_exit_code(3221225794)
assert "0xC0000142" in message
assert "DLL_INIT_FAILED" in message
def test_negative_signed_form_is_also_decoded(self):
# Python may hand back the signed form depending on how it was launched.
assert "0xC0000142" in describe_exit_code(-1073741502)
def test_unknown_code_degrades_to_the_raw_number(self):
assert describe_exit_code(1) == "Crawler exited with code 1"
# --------------------------------------------------------------------------
# Notes
# --------------------------------------------------------------------------
class TestNoteIngest:
@pytest.mark.asyncio
async def test_baseline_run_emits_no_new_note_events(self, db, tmp_path):
task = await _make_task(db)
run = await _make_run(db, task, started_at=1000)
_write_run_dir(tmp_path, [_note("n1"), _note("n2")], comments=[])
result = await ingest_run(db, run, task, tmp_path)
assert result.status == RUN_SUCCESS
assert result.is_baseline is True
assert result.new_notes == 2
# Everything is "new" on the first run; emitting that would be pure noise.
assert await _events(db, EVENT_NEW_NOTE) == []
assert len(list((await db.scalars(select(MonitorNote))).all())) == 2
@pytest.mark.asyncio
async def test_an_empty_run_does_not_establish_a_baseline(self, db, tmp_path):
"""A run that fetched nothing observed nothing, so it is not a baseline.
Otherwise the first crawl that actually works reports every work as
newly discovered.
"""
task = await _make_task(db)
empty_run = await _make_run(db, task, started_at=1000)
(tmp_path / "empty").mkdir(parents=True, exist_ok=True)
await ingest_run(db, empty_run, task, tmp_path / "empty")
real_run = await _make_run(db, task, started_at=2000)
result = await ingest_run(
db, real_run, task, _write_run_dir(tmp_path / "ok", [_note("n1")], comments=[])
)
assert result.is_baseline is True
assert await _events(db, EVENT_NEW_NOTE) == []
@pytest.mark.asyncio
async def test_second_run_reports_only_the_added_note(self, db, tmp_path):
task = await _make_task(db)
first_dir = _write_run_dir(tmp_path / "run1", [_note("n1")], comments=[])
run1 = await _make_run(db, task, started_at=1000)
await ingest_run(db, run1, task, first_dir)
second_dir = _write_run_dir(tmp_path / "run2", [_note("n1"), _note("n2")], comments=[])
run2 = await _make_run(db, task, started_at=2000)
result = await ingest_run(db, run2, task, second_dir)
assert result.is_baseline is False
assert result.new_notes == 1
events = await _events(db, EVENT_NEW_NOTE)
assert len(events) == 1
assert events[0].target_id == "n2"
assert events[0].run_id == run2.id
class TestMetricSnapshots:
@pytest.mark.asyncio
async def test_delta_event_emitted_when_like_count_changes(self, db, tmp_path):
task = await _make_task(db)
run1 = await _make_run(db, task, started_at=1000)
await ingest_run(db, run1, task, _write_run_dir(tmp_path / "r1", [_note("n1", "100")], comments=[]))
run2 = await _make_run(db, task, started_at=2000)
await ingest_run(db, run2, task, _write_run_dir(tmp_path / "r2", [_note("n1", "150")], comments=[]))
events = await _events(db, EVENT_METRIC_DELTA)
assert len(events) == 1
payload = json.loads(events[0].payload_json)
assert payload["deltas"]["liked_count"] == {"from": 100, "to": 150, "delta": 50}
@pytest.mark.asyncio
async def test_no_delta_when_nothing_changed(self, db, tmp_path):
task = await _make_task(db)
run1 = await _make_run(db, task, started_at=1000)
await ingest_run(db, run1, task, _write_run_dir(tmp_path / "r1", [_note("n1", "100")], comments=[]))
run2 = await _make_run(db, task, started_at=2000)
await ingest_run(db, run2, task, _write_run_dir(tmp_path / "r2", [_note("n1", "100")], comments=[]))
assert await _events(db, EVENT_METRIC_DELTA) == []
@pytest.mark.asyncio
async def test_unparseable_count_is_null_not_zero(self, db, tmp_path):
task = await _make_task(db)
run = await _make_run(db, task, started_at=1000)
await ingest_run(db, run, task, _write_run_dir(tmp_path, [_note("n1", "暂无")], comments=[]))
metric = await db.scalar(select(MonitorNoteMetric).where(MonitorNoteMetric.note_id == "n1"))
# Zero would forge a large negative delta on the next comparison.
assert metric.liked_count is None
assert metric.raw_liked_count == "暂无"
@pytest.mark.asyncio
async def test_no_delta_when_previous_value_was_unparseable(self, db, tmp_path):
task = await _make_task(db)
run1 = await _make_run(db, task, started_at=1000)
await ingest_run(db, run1, task, _write_run_dir(tmp_path / "r1", [_note("n1", "暂无")], comments=[]))
run2 = await _make_run(db, task, started_at=2000)
await ingest_run(db, run2, task, _write_run_dir(tmp_path / "r2", [_note("n1", "50")], comments=[]))
assert await _events(db, EVENT_METRIC_DELTA) == []
@pytest.mark.asyncio
async def test_metric_snapshot_survives_across_runs(self, db, tmp_path):
"""The crawler's own DB store overwrites metrics; ours must not."""
task = await _make_task(db)
for index, liked in enumerate(["100", "150", "300"]):
run = await _make_run(db, task, started_at=1000 * (index + 1))
await ingest_run(
db, run, task, _write_run_dir(tmp_path / f"r{index}", [_note("n1", liked)], comments=[])
)
snapshots = list(
(
await db.scalars(
select(MonitorNoteMetric)
.where(MonitorNoteMetric.note_id == "n1")
.order_by(MonitorNoteMetric.run_id)
)
).all()
)
assert [s.liked_count for s in snapshots] == [100, 150, 300]
# --------------------------------------------------------------------------
# Comments
# --------------------------------------------------------------------------
class TestCommentIngest:
@pytest.mark.asyncio
async def test_posted_vs_seen_split_by_create_time(self, db, tmp_path):
task = await _make_task(db)
# Baseline establishes the seen-set; no events on the first run.
run1 = await _make_run(db, task, started_at=1000)
await ingest_run(
db, run1, task,
_write_run_dir(tmp_path / "r1", [_note("n1")], comments=[_comment("c1", "n1", create_time=500)]),
)
assert await _events(db, EVENT_NEW_COMMENT_POSTED) == []
# c2 was published after run1 started -> genuinely new.
# c3 is old but only just surfaced in the top-N window -> seen, not posted.
run2 = await _make_run(db, task, started_at=2000)
await ingest_run(
db, run2, task,
_write_run_dir(
tmp_path / "r2",
[_note("n1")],
comments=[
_comment("c1", "n1", create_time=500),
_comment("c2", "n1", create_time=2500),
_comment("c3", "n1", create_time=100),
],
),
)
posted = await _events(db, EVENT_NEW_COMMENT_POSTED)
seen = await _events(db, EVENT_NEW_COMMENT_SEEN)
assert len(posted) == 1
assert json.loads(posted[0].payload_json)["comment_id"] == "c2"
assert len(seen) == 1
assert json.loads(seen[0].payload_json)["comment_id"] == "c3"
@pytest.mark.asyncio
async def test_comments_not_ingested_when_disabled(self, db, tmp_path):
task = await _make_task(db, enable_comments=False)
run = await _make_run(db, task, started_at=1000)
result = await ingest_run(
db, run, task,
_write_run_dir(tmp_path, [_note("n1")], comments=[_comment("c1", "n1", 500)]),
)
assert result.new_comments == 0
# --------------------------------------------------------------------------
# Failure handling
# --------------------------------------------------------------------------
class TestFailureHandling:
@pytest.mark.asyncio
async def test_nonzero_exit_is_a_failure(self, db, tmp_path):
task = await _make_task(db)
run = await _make_run(db, task, started_at=1000, exit_code=1)
_write_run_dir(tmp_path, [_note("n1")], comments=[])
result = await ingest_run(db, run, task, tmp_path)
assert result.status == RUN_FAILED
assert len(await _events(db, EVENT_RUN_FAILED)) == 1
# A crashed run must not touch the seen-set.
assert await db.scalar(select(MonitorNote.id)) is None
@pytest.mark.asyncio
async def test_zero_notes_with_exit_zero_is_a_suspected_auth_failure(self, db, tmp_path):
"""The silent-cookie-failure signature: exit 0 but nothing fetched.
A real bad-cookie run writes no output file at all, which is why the
exit code has to be checked before the files are.
"""
task = await _make_task(db)
run = await _make_run(db, task, started_at=1000, exit_code=0)
tmp_path.mkdir(parents=True, exist_ok=True)
result = await ingest_run(db, run, task, tmp_path)
assert result.status == RUN_PARTIAL
assert len(await _events(db, EVENT_AUTH_FAILURE)) == 1
assert await _events(db, EVENT_RUN_FAILED) == []
@pytest.mark.asyncio
async def test_no_data_is_not_blamed_on_the_cookie_when_a_sibling_succeeded(
self, db, tmp_path
):
"""A task that just worked proves the login is fine; do not cry wolf."""
healthy = await _make_task(db, name="healthy")
healthy_run = await _make_run(db, healthy, started_at=get_current_timestamp())
await ingest_run(
db, healthy_run, healthy,
_write_run_dir(tmp_path / "ok", [_note("n1")], comments=[]),
)
task = await _make_task(db, name="suspect")
run = await _make_run(db, task, started_at=get_current_timestamp())
(tmp_path / "empty").mkdir(parents=True, exist_ok=True)
result = await ingest_run(db, run, task, tmp_path / "empty")
assert result.status == RUN_PARTIAL
assert await _events(db, EVENT_NO_DATA) != []
assert await _events(db, EVENT_AUTH_FAILURE) == []
@pytest.mark.asyncio
async def test_empty_contents_file_is_also_an_auth_failure(self, db, tmp_path):
task = await _make_task(db)
run = await _make_run(db, task, started_at=1000, exit_code=0)
_write_run_dir(tmp_path, [], comments=[])
result = await ingest_run(db, run, task, tmp_path)
assert result.status == RUN_PARTIAL
assert len(await _events(db, EVENT_AUTH_FAILURE)) == 1
# --------------------------------------------------------------------------
# Idempotency
# --------------------------------------------------------------------------
class TestIdempotency:
@pytest.mark.asyncio
async def test_reingesting_the_same_data_adds_nothing(self, db, tmp_path):
task = await _make_task(db)
run_dir = _write_run_dir(
tmp_path, [_note("n1"), _note("n2")], comments=[_comment("c1", "n1", 500)]
)
run1 = await _make_run(db, task, started_at=1000)
await ingest_run(db, run1, task, run_dir)
notes_after_first = len(list((await db.scalars(select(MonitorNote))).all()))
# A retry of the same crawl content must not duplicate rows or events.
run2 = await _make_run(db, task, started_at=2000)
result = await ingest_run(db, run2, task, run_dir)
assert result.new_notes == 0
assert result.new_comments == 0
assert len(list((await db.scalars(select(MonitorNote))).all())) == notes_after_first