Files
MediaCrawler/tests/test_upstream.py
T
butubb e348de48d3
Deploy VitePress site to Pages / build (push) Canceled after 0s
Deploy VitePress site to Pages / Deploy (push) Canceled after 0s
feat(upstream): 上游更新检查——定期比上游、落后了推企业微信
本仓库在上游 MediaCrawler 之上加了一整层(见 UPSTREAM.md),可合并流程默认
「有人知道上游动了」。而部署是 git pull --ff-only,只从自己的 Gitea 拉——上游的
提交不主动 fetch 就永远看不见。拖着不合并的代价是复利的:越久越难合。

于是把「上游动了没有」变成一条会自己跑、会推企业微信的通知:

* api/monitor/upstream.py:git fetch <url> <branch> 到 FETCH_HEAD,用
  rev-list --count HEAD..FETCH_HEAD 算落后数、FETCH_HEAD..HEAD 算领先数。
  用 git 而非托管商 API,因为只有 git 知道共同祖先在哪——本仓库含有上游没有的
  提交,直接比 tip 会得出错误结论。增量 fetch 只传几个新提交,不会遇到
  UPSTREAM.md 里说的「大包必断」。
* 只 fetch 到 FETCH_HEAD:不配 remote、不写 refs/remotes、不碰索引与工作区,
  所以不打断正在跑的采集,也不和 deploy.sh 的 git pull 抢锁。
* 挂在调度器 tick 上(不是采集,所以不看 is_busy、不受活跃时段限制——定时检查
  放在半夜反而最合适),按 checked_at + 间隔 到期才跑;失败也写 checked_at,
  于是 GitHub 不通时是每间隔重试一次,而不是每个 tick 撞一次墙。
* 同一个 tip 只推一次(记 tip 而不是「推过没」),上游真又动了会再推。
* 两个接口:GET /monitor/upstream 只读缓存;POST /monitor/upstream/check 手动
  查一次且刻意不推通知——点按钮的人正看着结果。
* 默认关闭,间隔默认一天。

Dockerfile 显式装 git(python:slim 不带,而这是唯一的依赖);deploy.sh 顺带补上
一个真 bug 的提示:Dockerfile/requirements.txt 变了只 up -d 用的还是旧镜像。
2026-10-10 09:17:06 +08:00

421 lines
15 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_upstream.py
# GitHub: https://github.com/NanmiCoder
# Licensed under NON-COMMERCIAL LEARNING LICENSE 1.1
#
# 声明:本代码仅供学习和研究目的使用。使用者应遵守以下原则:
# 1. 不得用于任何商业用途。
# 2. 使用时应遵守目标平台的使用条款和robots.txt规则。
# 3. 不得进行大规模爬取或对平台造成运营干扰。
# 4. 应合理控制请求频率,避免给目标平台带来不必要的负担。
# 5. 不得用于任何非法或不当的用途。
#
# 详细许可条款请参阅项目根目录下的LICENSE文件。
# 使用本代码即表示您同意遵守上述原则和LICENSE中的所有条款。
"""上游更新检查:git 输出怎么解析,以及什么时候才推送。
所有会碰网络的路径都被替掉了 —— 测试里既没有上游仓库,也不该有。真正被测的是
解析、去重和调度到期这三件事,它们才是容易出错的部分。
"""
import subprocess
import httpx
import pytest
import pytest_asyncio
from api.main import app
from api.monitor import db as monitor_db
from api.monitor import notify
from api.monitor import scheduler as scheduler_module
from api.monitor import upstream
from api.monitor.models import SETTING_WECOM_WEBHOOK
from api.monitor.scheduler import MonitorScheduler
from api.monitor.settings import set_setting
from tools.time_util import get_current_timestamp
SEP = upstream._RECORD_SEPARATOR
HEAD_SHA = "a" * 40
TIP_SHA = "b" * 40
WEBHOOK = "https://qyapi.weixin.qq.com/cgi-bin/webhook/send?key=test"
@pytest_asyncio.fixture
async def db(tmp_path):
monitor_db.set_sqlite_path(tmp_path / "monitor.db")
await monitor_db.init_db()
yield monitor_db
await monitor_db.dispose_engine()
@pytest_asyncio.fixture
async def client(tmp_path):
monitor_db.set_sqlite_path(tmp_path / "monitor.db")
await monitor_db.init_db()
transport = httpx.ASGITransport(app=app)
async with httpx.AsyncClient(transport=transport, base_url="http://test") as http_client:
yield http_client
await monitor_db.dispose_engine()
def _signature(args: list[str]) -> str:
"""Map one git invocation onto a key a test can name.
``rev-parse`` and ``rev-list`` are each called more than once with different
arguments, so the subcommand alone is not enough to key on.
"""
if args[0] == "rev-parse":
return f"rev-parse {args[1]}"
if args[0] == "rev-list":
return f"rev-list {args[-1]}"
return args[0]
@pytest.fixture
def fake_git(monkeypatch):
"""Install canned git output. Returns the list of invocations made."""
def _install(responses: dict) -> list:
calls: list = []
def _run(args, timeout):
calls.append(args)
key = _signature(args)
if key not in responses:
raise AssertionError(f"测试没有为这条 git 调用准备输出:{args}")
return responses[key]
monkeypatch.setattr(upstream, "_run", _run)
return calls
return _install
def _log_line(sha: str, subject: str) -> str:
return f"{sha}{SEP}张三{SEP}2026-10-01{SEP}{subject}"
class TestCheckParsing:
def test_counts_commits_and_reads_the_list(self, fake_git):
calls = fake_git(
{
"rev-parse HEAD": (0, HEAD_SHA, ""),
"fetch": (0, "", ""),
"rev-parse FETCH_HEAD": (0, TIP_SHA, ""),
"rev-list HEAD..FETCH_HEAD": (0, "3", ""),
"rev-list FETCH_HEAD..HEAD": (0, "128", ""),
"log": (
0,
"\n".join(
[
_log_line("abc1234", "fix: 修了扫码过期"),
_log_line("def5678", "feat: 加了新平台"),
_log_line("9999999", "docs: 更新说明"),
]
),
"",
),
}
)
result = upstream._check_sync("https://example.invalid/repo.git", "main")
assert result.ok is True
assert result.behind == 3
# 领先数就是这一层的规模,合并时要一起保留,所以值得单独报出来。
assert result.ahead == 128
assert result.tip == TIP_SHA
assert result.head == HEAD_SHA
assert [commit.subject for commit in result.commits] == [
"fix: 修了扫码过期",
"feat: 加了新平台",
"docs: 更新说明",
]
assert result.commits[0].sha == "abc1234"
assert result.error == ""
def test_up_to_date_skips_reading_the_log(self, fake_git):
"""不落后时不该再去读提交列表 —— 那条 git log 没有意义。"""
calls = fake_git(
{
"rev-parse HEAD": (0, HEAD_SHA, ""),
"fetch": (0, "", ""),
"rev-parse FETCH_HEAD": (0, TIP_SHA, ""),
"rev-list HEAD..FETCH_HEAD": (0, "0", ""),
"rev-list FETCH_HEAD..HEAD": (0, "128", ""),
}
)
result = upstream._check_sync("https://example.invalid/repo.git", "main")
assert result.ok is True
assert result.behind == 0
assert result.commits == []
assert all(call[0] != "log" for call in calls)
def test_fetch_failure_is_reported_not_raised(self, fake_git):
"""GitHub 不通是常态,那也该是一条能显示出来的结论。"""
fake_git(
{
"rev-parse HEAD": (0, HEAD_SHA, ""),
"fetch": (128, "", "fatal: unable to access 'https://github.com/': 连接超时\n第二行"),
}
)
result = upstream._check_sync("https://example.invalid/repo.git", "main")
assert result.ok is False
assert result.head == HEAD_SHA
# 只留第一行 stderr:git 的报错常常跟一大段建议,塞进界面反而看不清。
assert "连接超时" in result.error
assert "第二行" not in result.error
def test_not_a_git_repository_is_reported(self, fake_git):
fake_git({"rev-parse HEAD": (128, "", "fatal: not a git repository (or any of the parent directories): .git")})
result = upstream._check_sync("https://example.invalid/repo.git", "main")
assert result.ok is False
assert "读取本地 HEAD 失败" in result.error
def test_missing_git_binary_is_reported(self, monkeypatch):
def _explode(args, timeout):
raise FileNotFoundError("git")
monkeypatch.setattr(upstream, "_git", _explode)
result = upstream._check_sync("https://example.invalid/repo.git", "main")
assert result.ok is False
assert "未找到 git" in result.error
def test_fetch_timeout_is_reported(self, monkeypatch):
def _explode(args, timeout):
raise subprocess.TimeoutExpired(cmd="git", timeout=timeout)
monkeypatch.setattr(upstream, "_git", _explode)
result = upstream._check_sync("https://example.invalid/repo.git", "main")
assert result.ok is False
assert "超时" in result.error
class TestNotification:
@pytest.fixture
def sent(self, monkeypatch) -> list:
messages: list = []
async def _fake_send(url, content):
messages.append(content)
return True, "发送成功"
monkeypatch.setattr(notify, "send_wecom", _fake_send)
return messages
@staticmethod
def _patch_check(monkeypatch, *, behind: int, tip: str, ahead: int = 0):
async def _fake_check(remote_url=upstream.DEFAULT_REMOTE_URL, branch=upstream.DEFAULT_BRANCH):
return upstream.CheckResult(
ok=True,
behind=behind,
ahead=ahead,
tip=tip,
head=HEAD_SHA,
commits=[upstream.Commit(sha="abc1234", author="张三", date="2026-10-01", subject="fix: 修了扫码过期")],
)
monkeypatch.setattr(upstream, "check", _fake_check)
@pytest.mark.asyncio
async def test_pushes_once_per_upstream_tip(self, db, monkeypatch, sent):
self._patch_check(monkeypatch, behind=2, tip=TIP_SHA, ahead=128)
async with monitor_db.get_session() as session:
await set_setting(session, SETTING_WECOM_WEBHOOK, WEBHOOK)
first = await upstream.run_check()
# 同一个 tip 再查一次:不该重复推。
second = await upstream.run_check()
assert len(sent) == 1
assert first["notified"] is True
assert "notified" not in second
assert "落后 `main` **2** 个提交" in sent[0]
assert "abc1234" in sent[0]
# 上游又动了:tip 变了就该再推一次。
self._patch_check(monkeypatch, behind=5, tip="c" * 40, ahead=128)
third = await upstream.run_check()
assert len(sent) == 2
assert third["notified"] is True
@pytest.mark.asyncio
async def test_state_is_persisted_for_the_ui(self, db, monkeypatch, sent):
self._patch_check(monkeypatch, behind=2, tip=TIP_SHA, ahead=128)
async with monitor_db.get_session() as session:
await set_setting(session, SETTING_WECOM_WEBHOOK, WEBHOOK)
await upstream.run_check()
async with monitor_db.get_session() as session:
state = await upstream.load_state(session)
assert state["behind"] == 2
assert state["ahead"] == 128
assert state["tip"] == TIP_SHA
assert state["branch"] == "main"
assert state["checked_at"] > 0
@pytest.mark.asyncio
async def test_nothing_new_does_not_push(self, db, monkeypatch, sent):
self._patch_check(monkeypatch, behind=0, tip=TIP_SHA)
async with monitor_db.get_session() as session:
await set_setting(session, SETTING_WECOM_WEBHOOK, WEBHOOK)
await upstream.run_check()
assert sent == []
@pytest.mark.asyncio
async def test_switch_off_does_not_push(self, db, monkeypatch, sent):
self._patch_check(monkeypatch, behind=2, tip=TIP_SHA)
async with monitor_db.get_session() as session:
await set_setting(session, SETTING_WECOM_WEBHOOK, WEBHOOK)
await set_setting(session, "system.upstream_notify", "false")
await upstream.run_check()
assert sent == []
@pytest.mark.asyncio
async def test_manual_check_can_skip_the_push(self, db, monkeypatch, sent):
"""手动点「立即检查」只看结果,不因为它把群消息推一遍。"""
self._patch_check(monkeypatch, behind=2, tip=TIP_SHA)
async with monitor_db.get_session() as session:
await set_setting(session, SETTING_WECOM_WEBHOOK, WEBHOOK)
await upstream.run_check(notify_when_new=False)
assert sent == []
@pytest.mark.asyncio
async def test_no_webhook_configured_is_not_an_error(self, db, monkeypatch, sent):
self._patch_check(monkeypatch, behind=2, tip=TIP_SHA)
result = await upstream.run_check()
assert sent == []
assert result["ok"] is True
assert result["behind"] == 2
class TestScheduler:
@pytest.fixture
def checks(self, monkeypatch) -> list:
calls: list = []
async def _fake_run_check(notify_when_new: bool = True):
calls.append(notify_when_new)
# 真实的 run_check 会把 checked_at 写进去,调度器的「到点了没有」
# 全靠这个字段,所以替身也必须写。
async with monitor_db.get_session() as session:
await upstream._save_state(
session,
{"checked_at": get_current_timestamp(), "ok": True, "behind": 0},
)
return {"ok": True, "behind": 0}
monkeypatch.setattr(upstream, "run_check", _fake_run_check)
return calls
@pytest.mark.asyncio
async def test_disabled_never_checks(self, db, checks):
await MonitorScheduler()._maybe_check_upstream()
assert checks == []
@pytest.mark.asyncio
async def test_checks_when_due_and_then_waits_out_the_interval(self, db, checks):
async with monitor_db.get_session() as session:
await set_setting(session, "system.upstream_check_enabled", "true")
scheduler = MonitorScheduler()
await scheduler._maybe_check_upstream()
assert checks == [True]
# 刚查过:间隔(默认一天)没到就不该再查。
await scheduler._maybe_check_upstream()
assert checks == [True]
@pytest.mark.asyncio
async def test_a_stale_timestamp_is_due_again(self, db, checks):
async with monitor_db.get_session() as session:
await set_setting(session, "system.upstream_check_enabled", "true")
await set_setting(session, "system.upstream_check_interval_minutes", "30")
await upstream._save_state(
session,
# 差一分钟就到期,用来卡住边界:31 分钟前那次已经算过期。
{"checked_at": get_current_timestamp() - 31 * 60_000, "ok": True, "behind": 0},
)
await MonitorScheduler()._maybe_check_upstream()
assert checks == [True]
class TestEndpoint:
@pytest.mark.asyncio
async def test_status_is_empty_before_the_first_check(self, client):
response = await client.get("/api/monitor/upstream")
assert response.status_code == 200
assert response.json() == {}
@pytest.mark.asyncio
async def test_manual_check_runs_and_is_readable_back(self, client, monkeypatch):
async def _fake_check(remote_url=upstream.DEFAULT_REMOTE_URL, branch=upstream.DEFAULT_BRANCH):
return upstream.CheckResult(ok=True, behind=1, tip=TIP_SHA, head=HEAD_SHA)
monkeypatch.setattr(upstream, "check", _fake_check)
checked = (await client.post("/api/monitor/upstream/check")).json()
assert checked["behind"] == 1
# 结果落库,随后的 GET 读的是同一份缓存(而不是再 fetch 一次)。
cached = (await client.get("/api/monitor/upstream")).json()
assert cached["behind"] == 1
assert cached["tip"] == TIP_SHA
@pytest.mark.asyncio
async def test_manual_check_does_not_push(self, client, monkeypatch):
"""点按钮的人正看着结果,不该再给自己推一条群消息。"""
sent: list = []
async def _fake_send(url, content):
sent.append(content)
return True, "发送成功"
async def _fake_check(remote_url=upstream.DEFAULT_REMOTE_URL, branch=upstream.DEFAULT_BRANCH):
return upstream.CheckResult(ok=True, behind=1, tip=TIP_SHA, head=HEAD_SHA)
monkeypatch.setattr(notify, "send_wecom", _fake_send)
monkeypatch.setattr(upstream, "check", _fake_check)
async with monitor_db.get_session() as session:
await set_setting(session, SETTING_WECOM_WEBHOOK, WEBHOOK)
body = (await client.post("/api/monitor/upstream/check")).json()
assert sent == []
assert body["behind"] == 1
# 没推过的那批提交留给下一次定时检查,所以这里不该记成已推送。
async with monitor_db.get_session() as session:
state = await upstream.load_state(session)
assert "notified" not in state