From e348de48d3ac1eafa99fe2f30753908fa5c21217 Mon Sep 17 00:00:00 2001
From: butubb <1422726308@qq.com>
Date: Sat, 10 Oct 2026 09:17:06 +0800
Subject: [PATCH] =?UTF-8?q?feat(upstream):=20=E4=B8=8A=E6=B8=B8=E6=9B=B4?=
=?UTF-8?q?=E6=96=B0=E6=A3=80=E6=9F=A5=E2=80=94=E2=80=94=E5=AE=9A=E6=9C=9F?=
=?UTF-8?q?=E6=AF=94=E4=B8=8A=E6=B8=B8=E3=80=81=E8=90=BD=E5=90=8E=E4=BA=86?=
=?UTF-8?q?=E6=8E=A8=E4=BC=81=E4=B8=9A=E5=BE=AE=E4=BF=A1?=
MIME-Version: 1.0
Content-Type: text/plain; charset=UTF-8
Content-Transfer-Encoding: 8bit
本仓库在上游 MediaCrawler 之上加了一整层(见 UPSTREAM.md),可合并流程默认
「有人知道上游动了」。而部署是 git pull --ff-only,只从自己的 Gitea 拉——上游的
提交不主动 fetch 就永远看不见。拖着不合并的代价是复利的:越久越难合。
于是把「上游动了没有」变成一条会自己跑、会推企业微信的通知:
* api/monitor/upstream.py:git fetch ` 显示。
@@ -89,6 +94,16 @@ Dockerfile / .dockerignore / docker-compose.yml 服务器部署用
## 三、上游更新时怎么操作
+### 先让机器替你盯着
+
+「上游更新检查」(`api/monitor/upstream.py`,开关在 WebUI 的**系统设置 → 上游更新**)会按
+间隔 `git fetch` 上游、算出落后几个提交,有更新就推企业微信。它是这份文档的自动化版:
+没有它,「上游动了」这件事只取决于谁偶尔想起来去 fetch 一次。
+
+两个细节决定了它为什么是安全的:它只 fetch 到 `FETCH_HEAD`,**不写工作区、不建 remote、不碰
+`refs/remotes`**,所以和正在跑的采集、和下面的 `git pull` 都不冲突;默认**关闭**,因为要联网,
+且需要镜像里有 git。
+
### 日常流程
```bash
diff --git a/api/monitor/app_settings.py b/api/monitor/app_settings.py
index b98fabc..539cf91 100644
--- a/api/monitor/app_settings.py
+++ b/api/monitor/app_settings.py
@@ -53,6 +53,7 @@ from .settings import (
set_setting,
system_key,
)
+from .upstream import DEFAULT_BRANCH, DEFAULT_REMOTE_URL
SCOPE_PLATFORM = "platform"
SCOPE_SYSTEM = "system"
@@ -216,6 +217,61 @@ SETTING_SPECS: List[SettingSpec] = [
),
default=False,
),
+ # --- 上游更新检查 -------------------------------------------------------
+ # 这几项不作用于采集,所以都标 affects_new_runs=False:改动它们不需要等下一轮,
+ # 也不影响采集命令的拼装。
+ SettingSpec(
+ name="upstream_check_enabled",
+ scope=SCOPE_SYSTEM,
+ type=TYPE_BOOL,
+ label="检查上游仓库更新",
+ help=(
+ "定期 fetch 上游仓库,看看它有没有新提交,并在有更新时推送通知。"
+ "本仓库在上游之上加了一整层(见 UPSTREAM.md),不定期看一眼就会越拖越难合并。"
+ ),
+ default=False,
+ affects_new_runs=False,
+ ),
+ SettingSpec(
+ name="upstream_check_interval_minutes",
+ scope=SCOPE_SYSTEM,
+ type=TYPE_INT,
+ label="上游检查间隔(分钟)",
+ help="默认 1440 分钟(每天一次)。检查只是 fetch,不需要太频繁。",
+ default=1440,
+ minimum=30,
+ maximum=10080,
+ affects_new_runs=False,
+ ),
+ SettingSpec(
+ name="upstream_remote_url",
+ scope=SCOPE_SYSTEM,
+ type=TYPE_STR,
+ label="上游仓库地址",
+ help=(
+ "默认是 GitHub 上的上游。国内直连 GitHub 不稳时改成 gitcode 镜像"
+ "(见 UPSTREAM.md),或任意能访问到上游的地址。"
+ ),
+ default=DEFAULT_REMOTE_URL,
+ affects_new_runs=False,
+ ),
+ SettingSpec(
+ name="upstream_branch",
+ scope=SCOPE_SYSTEM,
+ type=TYPE_STR,
+ label="上游分支",
+ default=DEFAULT_BRANCH,
+ affects_new_runs=False,
+ ),
+ SettingSpec(
+ name="upstream_notify",
+ scope=SCOPE_SYSTEM,
+ type=TYPE_BOOL,
+ label="上游有更新时推送通知",
+ help="只在出现此前没推过的上游提交时发一条,同一个更新不会反复推。",
+ default=True,
+ affects_new_runs=False,
+ ),
]
SPECS_BY_NAME = {spec.name: spec for spec in SETTING_SPECS}
diff --git a/api/monitor/models.py b/api/monitor/models.py
index 596d52d..091362c 100644
--- a/api/monitor/models.py
+++ b/api/monitor/models.py
@@ -357,6 +357,16 @@ SETTING_AUTH_PASSWORD_UPDATED_AT = "auth_password_updated_at"
# them. Key builders live in settings.py.
SETTING_WECOM_WEBHOOK = "system.wecom_webhook"
+# 上游更新检查的两条状态。都不是给用户编辑的设置项,所以不在 app_settings 的注册表里
+# (那张表只列可编辑项,因此也不会被设置接口读出来)。
+#
+# 最近一次检查的结果整体存成一条 JSON:它总是被整体读写,拆成多个 key 只会带来
+# 另一半没写完的不一致。
+SETTING_UPSTREAM_STATE = "system.upstream_check_state"
+# 已经推送过通知的那个上游 tip。换 tip 才再推 —— 否则每个检查周期都会把同样的
+# 更新推一遍,直到有人去合并为止;而上游真又动了的时候应该再推一次。
+SETTING_UPSTREAM_NOTIFIED_TIP = "system.upstream_notified_tip"
+
# Pre-namespacing keys, kept only so the startup migration can find and move
# them. Nothing should read these directly.
LEGACY_SETTING_KEY_RENAMES = {
diff --git a/api/monitor/scheduler.py b/api/monitor/scheduler.py
index cf6c63b..6de3ddd 100644
--- a/api/monitor/scheduler.py
+++ b/api/monitor/scheduler.py
@@ -32,6 +32,11 @@ Two families of schedule, and the difference matters:
calendar, so a run that starts late does not drag every later run with it.
The arithmetic for both lives in schedule.py.
+
+The loop also carries the 上游更新检查: it is not a crawl, so it shares none of
+the rules above (no subprocess, no active-hours gate) -- see
+``_maybe_check_upstream``. It rides this loop rather than getting a thread of its
+own because it is one HTTP-shaped fetch per day.
"""
import asyncio
@@ -44,7 +49,7 @@ from sqlalchemy import select
from tools.time_util import get_current_timestamp
from ..services import crawler_manager
-from . import app_settings, schedule
+from . import app_settings, schedule, upstream
from .db import get_session
from .models import MonitorRun, MonitorTask, RUN_INTERRUPTED, RUN_RUNNING
from .runner import execute_task
@@ -92,8 +97,49 @@ class MonitorScheduler:
await self.tick()
except Exception as exc: # pragma: no cover - keep the loop alive
print(f"[monitor.scheduler] tick failed: {exc}")
+ # 独立于采集任务,因此单独一段 try:上游检查失败不该影响采集调度,
+ # 反过来也一样。
+ try:
+ await self._maybe_check_upstream()
+ except Exception as exc: # pragma: no cover - keep the loop alive
+ print(f"[monitor.scheduler] upstream check failed: {exc}")
await asyncio.sleep(POLL_INTERVAL_SECONDS)
+ async def _maybe_check_upstream(self) -> None:
+ """到点就 fetch 一次上游仓库,看它有没有新提交。
+
+ 与采集任务的三条规则都不同,各有理由:它不碰浏览器、也不占采集子进程,
+ 所以不看 ``is_busy``;它只发一个 git 请求,没有被平台风控的风险,所以也不
+ 受活跃时段限制 —— 定时检查放在半夜反而是最合适的。
+ """
+ async with get_session() as session:
+ if not await app_settings.get_value(
+ session, "upstream_check_enabled", fallback=False
+ ):
+ return
+ interval_minutes = int(
+ await app_settings.get_value(
+ session, "upstream_check_interval_minutes", fallback=1440
+ )
+ )
+ state = await upstream.load_state(session)
+
+ checked_at = int(state.get("checked_at") or 0)
+ now = get_current_timestamp()
+ # 失败也会写 checked_at,所以不通的时候同样是每个间隔重试一次,
+ # 而不是每个 tick(20 秒)都去撞一次墙。
+ if checked_at and now - checked_at < max(1, interval_minutes) * 60_000:
+ return
+
+ result = await upstream.run_check()
+ if result.get("behind"):
+ print(
+ f"[monitor.scheduler] 上游 {result.get('branch')} 领先 "
+ f"{result['behind']} 个提交"
+ )
+ elif not result.get("ok"):
+ print(f"[monitor.scheduler] 上游检查失败:{result.get('error')}")
+
async def recover(self) -> None:
"""Clean up state left behind by a server restart.
diff --git a/api/monitor/upstream.py b/api/monitor/upstream.py
new file mode 100644
index 0000000..771abaf
--- /dev/null
+++ b/api/monitor/upstream.py
@@ -0,0 +1,355 @@
+# -*- coding: utf-8 -*-
+# Copyright (c) 2025 relakkes@gmail.com
+#
+# This file is part of MediaCrawler project.
+# Repository: https://github.com/NanmiCoder/MediaCrawler/blob/main/api/monitor/upstream.py
+# GitHub: https://github.com/NanmiCoder
+# Licensed under NON-COMMERCIAL LEARNING LICENSE 1.1
+#
+# 声明:本代码仅供学习和研究目的使用。使用者应遵守以下原则:
+# 1. 不得用于任何商业用途。
+# 2. 使用时应遵守目标平台的使用条款和robots.txt规则。
+# 3. 不得进行大规模爬取或对平台造成运营干扰。
+# 4. 应合理控制请求频率,避免给目标平台带来不必要的负担。
+# 5. 不得用于任何非法或不当的用途。
+#
+# 详细许可条款请参阅项目根目录下的LICENSE文件。
+# 使用本代码即表示您同意遵守上述原则和LICENSE中的所有条款。
+
+"""上游仓库更新检查。
+
+本仓库在上游(NanmiCoder/MediaCrawler)之上加了一整层监控/鉴权/多平台面板,
+合并流程写在 UPSTREAM.md 里。但那份流程默认**有人知道上游动了**——而部署脚本是
+`git pull --ff-only`,只从我们自己的 Gitea 拉,上游的提交不主动去 fetch 就永远
+看不见。拖着不合并的代价是复利的:越久越难合,最后只能放弃。这个模块把「上游动
+了没有」变成一条可定时、会推到企业微信的通知。
+
+三处刻意的取舍:
+
+* **用 git 而不是托管商的 HTTP API。** 只有 git 算得出「落后几个提交」:托管商
+ API 能告诉你上游 tip 是什么,但它不知道我们与上游的共同祖先在哪,而分歧点恰恰
+ 是真正要合的东西。本仓库还含有上游没有的提交,直接比 tip 会得出错误的结论。
+* **按 URL fetch 到 FETCH_HEAD,不配置 remote、不写 refs/remotes。** 服务器上的
+ checkout 是从 Gitea 克隆的,本来就没有 upstream 这个 remote;用 URL 直取就不必
+ 先去改它的 git 配置。顺带也避免往别人的部署里塞一个 remote。
+* **只读不写工作区。** fetch 只落对象和 FETCH_HEAD,不碰索引与工作区,所以不会打断
+ 正在跑的采集,也不会和 `./deploy.sh` 的 git pull 抢锁。
+
+依赖一个外部命令:**git**。本机开发环境一定有;容器里是 Dockerfile 显式装的
+(python:3.11-slim 默认不带)。
+"""
+
+import asyncio
+import json
+import os
+import subprocess
+from dataclasses import dataclass, field
+from pathlib import Path
+from typing import Any, Dict, List, Tuple
+
+from sqlalchemy.ext.asyncio import AsyncSession
+
+from tools.time_util import get_current_timestamp
+
+from .db import get_session
+from .models import SETTING_UPSTREAM_NOTIFIED_TIP, SETTING_UPSTREAM_STATE
+from .settings import get_setting, set_setting
+
+PROJECT_ROOT = Path(__file__).parent.parent.parent
+
+# 默认就是本仓库跟踪的那个上游。国内直连 GitHub 不稳时改成 gitcode 镜像即可
+# (见 UPSTREAM.md「直连 GitHub 不通时」)。
+DEFAULT_REMOTE_URL = "https://github.com/NanmiCoder/MediaCrawler.git"
+DEFAULT_BRANCH = "main"
+
+# fetch 要走网络,给宽松些;其余全是本地命令,慢到这个程度只能说明仓库坏了。
+FETCH_TIMEOUT_SECONDS = 120
+LOCAL_TIMEOUT_SECONDS = 20
+
+# 通知里最多列几条提交。要传达的是「该动手了」,不是把 changelog 搬到群里。
+MAX_LISTED_COMMITS = 10
+# 状态里留几条给前端展示。比通知多留一些,界面上能看到更完整的列表。
+MAX_STORED_COMMITS = 30
+
+# git log 用 Unit Separator 分隔字段:它不可能出现在提交信息里,比制表符安全。
+_RECORD_SEPARATOR = "\x1f"
+_LOG_FORMAT = (
+ f"%h{_RECORD_SEPARATOR}%an{_RECORD_SEPARATOR}%ad{_RECORD_SEPARATOR}%s"
+)
+
+
+class GitError(RuntimeError):
+ """git 不可用,或某条 git 命令失败了。"""
+
+
+@dataclass(frozen=True)
+class Commit:
+ """一条上游提交,只留通知/展示需要的四个字段。"""
+
+ sha: str
+ author: str
+ date: str
+ subject: str
+
+
+@dataclass
+class CheckResult:
+ """一次检查的结论。失败也是一种结论,用 ``ok``/``error`` 表达而不是抛异常。"""
+
+ ok: bool
+ # HEAD..FETCH_HEAD:上游有而我们没有的提交数 —— 要合的就是这些。
+ behind: int = 0
+ # FETCH_HEAD..HEAD:我们有自己的提交数 —— 也就是这一层的规模。
+ ahead: int = 0
+ tip: str = ""
+ head: str = ""
+ commits: List[Commit] = field(default_factory=list)
+ error: str = ""
+
+ def as_dict(self) -> Dict[str, Any]:
+ return {
+ "ok": self.ok,
+ "behind": self.behind,
+ "ahead": self.ahead,
+ "tip": self.tip,
+ "head": self.head,
+ "commits": [commit.__dict__ for commit in self.commits],
+ "error": self.error,
+ }
+
+
+def _git(args: List[str], timeout: int) -> subprocess.CompletedProcess:
+ env = dict(os.environ)
+ # 远端要凭据时(地址写成了私有仓库),git 会停下来问密码,而这里没有终端可问,
+ # 于是挂到超时。关掉一切交互,让它立刻失败。
+ env["GIT_TERMINAL_PROMPT"] = "0"
+ env["GIT_ASKPASS"] = ""
+ env["SSH_ASKPASS"] = ""
+ # 容器里 uid 1000 没有 passwd 项,git 找不到 HOME 会抱怨。给一个存在且可写的。
+ env.setdefault("HOME", "/tmp")
+
+ return subprocess.run(
+ [
+ "git",
+ "-C",
+ str(PROJECT_ROOT),
+ # 只对自己这个 checkout 放行所有权检查。容器里 uid 一般与属主一致,
+ # 但 bind mount 的属主未必,一旦不一致 git 会直接拒绝干任何活。
+ "-c",
+ f"safe.directory={PROJECT_ROOT}",
+ # 忽略任何全局凭据助手:这是个只读的公开仓库,不该去翻钥匙串。
+ "-c",
+ "credential.helper=",
+ *args,
+ ],
+ capture_output=True,
+ text=True,
+ encoding="utf-8",
+ errors="replace",
+ timeout=timeout,
+ env=env,
+ )
+
+
+def _run(args: List[str], timeout: int) -> Tuple[int, str, str]:
+ """跑一条 git 命令,返回 (returncode, stdout, stderr)。
+
+ 只把「跑不起来」当异常;命令返回非零是正常结果,交给调用方处理。
+ """
+ try:
+ proc = _git(args, timeout)
+ except FileNotFoundError as exc:
+ raise GitError("未找到 git 命令,请先安装 git") from exc
+ except subprocess.TimeoutExpired as exc:
+ raise GitError(f"git {args[0]} 超时({timeout} 秒)") from exc
+ return proc.returncode, (proc.stdout or "").strip(), (proc.stderr or "").strip()
+
+
+def _require(args: List[str], timeout: int, what: str) -> str:
+ code, out, err = _run(args, timeout)
+ if code != 0:
+ # git 的报错通常是多行的,只留第一行;完整输出塞进日志反而更难读。
+ detail = err.splitlines()[0].strip() if err else "未知错误"
+ raise GitError(f"{what}:{detail}")
+ return out
+
+
+def _to_int(raw: str) -> int:
+ try:
+ return int(raw)
+ except (TypeError, ValueError):
+ return 0
+
+
+def _parse_log(raw: str) -> List[Commit]:
+ commits: List[Commit] = []
+ for line in raw.splitlines():
+ parts = line.split(_RECORD_SEPARATOR)
+ if len(parts) != 4:
+ # 格式不对就跳过这一条:一条读不出来的提交不该让整次检查失败。
+ continue
+ sha, author, date, subject = parts
+ commits.append(Commit(sha=sha, author=author, date=date, subject=subject))
+ return commits
+
+
+def _check_sync(remote_url: str, branch: str) -> CheckResult:
+ """阻塞实现,异步包装见 :func:`check`。"""
+ # 在 try 之前绑定:后面的失败结果也带上它 —— 「检查失败」时当前跑的是哪个
+ # 提交,正是排查时第一个想知道的。
+ head = ""
+ try:
+ head = _require(["rev-parse", "HEAD"], LOCAL_TIMEOUT_SECONDS, "读取本地 HEAD 失败")
+ # 增量 fetch:对象本地基本都已经有了,所以正常情况下只传几个新提交,
+ # 不会遇到 UPSTREAM.md 里说的「大包必断」。
+ _require(
+ ["fetch", "--no-tags", remote_url, branch],
+ FETCH_TIMEOUT_SECONDS,
+ "从上游 fetch 失败",
+ )
+ tip = _require(["rev-parse", "FETCH_HEAD"], LOCAL_TIMEOUT_SECONDS, "读不到 FETCH_HEAD")
+ behind = _to_int(
+ _require(
+ ["rev-list", "--count", "HEAD..FETCH_HEAD"],
+ LOCAL_TIMEOUT_SECONDS,
+ "统计落后提交数失败",
+ )
+ )
+ ahead = _to_int(
+ _require(
+ ["rev-list", "--count", "FETCH_HEAD..HEAD"],
+ LOCAL_TIMEOUT_SECONDS,
+ "统计领先提交数失败",
+ )
+ )
+ # 只在确实落后时才读提交列表:已经是最新时这条 git log 毫无意义。
+ raw_log = ""
+ if behind:
+ raw_log = _require(
+ [
+ "log",
+ f"--max-count={MAX_STORED_COMMITS}",
+ "--date=short",
+ f"--format={_LOG_FORMAT}",
+ "HEAD..FETCH_HEAD",
+ ],
+ LOCAL_TIMEOUT_SECONDS,
+ "读取新提交列表失败",
+ )
+ except GitError as exc:
+ return CheckResult(ok=False, head=head, error=str(exc))
+
+ return CheckResult(
+ ok=True,
+ behind=behind,
+ ahead=ahead,
+ tip=tip,
+ head=head,
+ commits=_parse_log(raw_log),
+ )
+
+
+async def check(
+ remote_url: str = DEFAULT_REMOTE_URL, branch: str = DEFAULT_BRANCH
+) -> CheckResult:
+ """跑一次检查。网络与子进程都丢进线程,事件循环不被阻塞。
+
+ 不抛异常:上游不通是常态(尤其是直连 GitHub),那也是一种要记录下来的结果。
+ """
+ return await asyncio.to_thread(_check_sync, remote_url, branch)
+
+
+def build_message(result: CheckResult, branch: str) -> str:
+ """把一次「上游有新提交」的结果写成一条企业微信 markdown。"""
+ lines = [
+ "**🔔 上游 MediaCrawler 有更新**",
+ f"> 当前部署落后 `{branch}` **{result.behind}** 个提交",
+ ]
+ if result.ahead:
+ lines.append(f"> (本仓库另有 {result.ahead} 个自己的提交,合并时注意保留)")
+ for commit in result.commits[:MAX_LISTED_COMMITS]:
+ lines.append(f"> `{commit.sha}` {commit.subject}")
+ # 落后数可能大于列出来的条数:状态里留的提交本身也是截断的(30 条),
+ # 所以这里比的是总数,不是 len(commits)。
+ if result.behind > MAX_LISTED_COMMITS:
+ lines.append(f"> …等共 {result.behind} 个提交")
+ lines.append("> 合并步骤见仓库根目录 `UPSTREAM.md`")
+ return "\n".join(lines)
+
+
+async def load_state(session: AsyncSession) -> Dict[str, Any]:
+ """最近一次检查的结果,从设置里读回来。没查过时是空字典。"""
+ raw = await get_setting(session, SETTING_UPSTREAM_STATE)
+ if not raw:
+ return {}
+ try:
+ state = json.loads(raw)
+ except json.JSONDecodeError:
+ # 手改坏了的行不该让接口 500,当作「没查过」即可。
+ return {}
+ return state if isinstance(state, dict) else {}
+
+
+async def _save_state(session: AsyncSession, state: Dict[str, Any]) -> None:
+ # 整体读写,所以存成一条 JSON:拆成多个 key 只会带来写到一半的不一致。
+ await set_setting(session, SETTING_UPSTREAM_STATE, json.dumps(state, ensure_ascii=False))
+
+
+async def run_check(notify_when_new: bool = True) -> Dict[str, Any]:
+ """检查一次,落库,必要时推送。返回值可直接交给前端。
+
+ 分三段各自的数据库会话:fetch 最长可能跑满两分钟,占着一个连接不合适 ——
+ 理由与 runner.py 的分段完全相同。
+ """
+ # 延迟导入:app_settings 在模块级 import 本模块(为了那个默认地址常量),
+ # 模块级反向 import 会成环。
+ from . import app_settings, notify
+
+ async with get_session() as session:
+ remote_url = str(
+ await app_settings.get_value(
+ session, "upstream_remote_url", fallback=DEFAULT_REMOTE_URL
+ )
+ or DEFAULT_REMOTE_URL
+ )
+ branch = str(
+ await app_settings.get_value(session, "upstream_branch", fallback=DEFAULT_BRANCH)
+ or DEFAULT_BRANCH
+ )
+ notify_enabled = bool(
+ await app_settings.get_value(session, "upstream_notify", fallback=True)
+ )
+ notified_tip = (await get_setting(session, SETTING_UPSTREAM_NOTIFIED_TIP)) or ""
+ webhook_url = await notify.get_webhook_url(session)
+
+ result = await check(remote_url, branch)
+
+ payload: Dict[str, Any] = {
+ "checked_at": get_current_timestamp(),
+ "remote_url": remote_url,
+ "branch": branch,
+ **result.as_dict(),
+ }
+
+ async with get_session() as session:
+ await _save_state(session, payload)
+
+ has_update = result.ok and result.behind > 0 and bool(result.tip)
+ # 同一个 tip 只推一次:否则每过一个检查周期就把同样的更新推到群里,
+ # 直到有人去合为止。上游真又动了(tip 变了)时应该再推。
+ if (
+ notify_when_new
+ and notify_enabled
+ and has_update
+ and webhook_url
+ and result.tip != notified_tip
+ ):
+ ok, detail = await notify.send_wecom(webhook_url, build_message(result, branch))
+ if ok:
+ await set_setting(session, SETTING_UPSTREAM_NOTIFIED_TIP, result.tip)
+ payload["notified"] = True
+ else:
+ # 推送失败不该抹掉检查结果 —— 界面上仍然能看到「落后几个提交」。
+ payload["notify_error"] = detail
+
+ return payload
diff --git a/api/routers/monitor.py b/api/routers/monitor.py
index 08c671e..26c9e51 100644
--- a/api/routers/monitor.py
+++ b/api/routers/monitor.py
@@ -24,7 +24,7 @@ from typing import Any, Dict, List, Optional
from fastapi import APIRouter, HTTPException, Query, Response
from fastapi.responses import FileResponse
-from ..monitor import covers, notify, qrlogin, report, service
+from ..monitor import covers, notify, qrlogin, report, service, upstream
from ..monitor.db import get_session
from ..monitor.platforms import PLATFORM_XHS
from ..monitor.settings import (
@@ -596,3 +596,32 @@ async def test_webhook(payload: WebhookTestPayload):
if not ok:
raise HTTPException(status_code=400, detail=detail)
return {"message": detail}
+
+
+# ---------------------------------------------------------------------------
+# 上游更新
+# ---------------------------------------------------------------------------
+
+
+@router.get("/upstream")
+async def get_upstream_status():
+ """最近一次上游检查的结果。
+
+ 只读缓存,不触发检查:fetch 要走网络、最长两分钟,不该由一个 GET 顺手发起。
+ 没有查过时返回空对象,前端据此显示「尚未检查」。
+ """
+ async with get_session() as session:
+ return await upstream.load_state(session)
+
+
+@router.post("/upstream/check")
+async def run_upstream_check():
+ """立刻检查一次上游仓库,并把结果写回缓存。
+
+ 即使「定期检查」开关是关的也照查 —— 手动点这一次的意义正在于此。这里会一直
+ 等到 fetch 结束(前端给这条请求单独放长了超时),因为结果就是要给人看的。
+
+ ``notify_when_new=False``:点这个按钮的人正看着结果,没必要再给自己推一条群消息。
+ 没有推过的那批提交会留给下一次「定时检查」推 —— 推送状态记的是 tip,不是「推过没」。
+ """
+ return await upstream.run_check(notify_when_new=False)
diff --git a/deploy.sh b/deploy.sh
index d5ab63c..8b1bca2 100755
--- a/deploy.sh
+++ b/deploy.sh
@@ -27,6 +27,14 @@ else
git --no-pager log --oneline "$before..$after" | sed 's/^/ /'
fi
+# 镜像层(依赖)改动只能靠重建,而这一步不是自动的:Dockerfile 或 requirements.txt 变了,
+# 下面那句 `docker compose up -d --force-recreate` 用的是旧镜像,改动根本不会生效。
+# 至少要说出来,否则现象是「代码明明更新了,功能却报缺依赖」。
+if [ "$before" != "$after" ] && ! git diff --quiet "$before" "$after" -- Dockerfile requirements.txt; then
+ echo "!! Dockerfile / requirements.txt 有改动,需要重建镜像后重跑本脚本:"
+ echo " docker compose build"
+fi
+
# 前端重建的两种情况:产物根本不存在(首次部署),或 webui/ 有改动。
if [ ! -f api/webui/index.html ]; then
need_build=1
diff --git a/docs/监控功能使用说明.md b/docs/监控功能使用说明.md
index dfec3fa..7437055 100644
--- a/docs/监控功能使用说明.md
+++ b/docs/监控功能使用说明.md
@@ -297,7 +297,7 @@ set MC_PASSWORD=我的新密码 # Windows cmd
| 位置 | 范围 | 内容 |
|---|---|---|
| 左侧导航「设置」 | **按平台** | 登录 Cookie、采集策略、代理 |
-| 右上角「系统设置」 | **全局** | 通知、活跃时段、账号安全 |
+| 右上角「系统设置」 | **全局** | 通知、活跃时段、上游更新、账号安全 |
**这不是随便分的**:企业微信只有一个群、调度器只有一套时段规则、密码只有一份 ——
把它们放进"小红书专属"的页面里,会让人以为它们是按平台存的。
@@ -308,6 +308,54 @@ set MC_PASSWORD=我的新密码 # Windows cmd
---
+## 二·十一、上游更新检查
+
+本仓库在 [NanmiCoder/MediaCrawler](https://github.com/NanmiCoder/MediaCrawler) 之上加了一整层
+(监控 / 鉴权 / 多平台面板),差异管理与合并流程在根目录 `UPSTREAM.md` 里。但那份流程有个
+隐含前提:**得有人知道上游动了**。部署脚本只从我们自己的 Gitea `git pull`,上游的提交不主动
+去 fetch 就永远看不见 —— 拖着不合并的代价是复利的,越久越难合。
+
+这一项就是替你定时去 fetch 的:按间隔(默认每天一次)拉一次上游,算出「当前部署落后几个
+提交」,有更新就推一条企业微信,并把结果与提交列表显示在**右上角「系统设置」→「上游更新」**。
+
+### 配置
+
+| 项 | 默认 | 说明 |
+|---|---|---|
+| 检查上游仓库更新 | **关** | 总开关。默认关:它要联网 fetch,且需要容器里有 git(见下) |
+| 上游检查间隔(分钟) | 1440 | 每天一次。最小 30 分钟 |
+| 上游仓库地址 | GitHub 上游 | 国内直连 GitHub 不稳时改成 gitcode 镜像,见 `UPSTREAM.md` |
+| 上游分支 | `main` | |
+| 上游有更新时推送通知 | 开 | 只在出现**此前没推过**的提交时发一条,同一个更新不会反复推 |
+
+### 几个刻意的行为
+
+- **只读,不写工作区**:只 `git fetch <地址> <分支>` 到 `FETCH_HEAD` —— 不建 remote、不写
+ `refs/remotes`、不碰索引与工作区。所以它不会打断正在跑的采集,也不会和 `./deploy.sh`
+ 的 `git pull` 抢锁。
+- **不受活跃时段限制**:活跃时段是给采集定的(避免半夜去抓平台)。检查只是 fetch 一个公开
+ 仓库,半夜跑反而更合适。
+- **失败也是一种结果**:上游不通(尤其直连 GitHub)很常见。界面会显示失败原因与上次检查
+ 时间,失败不推送,也**不会**因此改变下一次检查的时间 —— 每个间隔重试一次,而不是每个
+ 调度 tick(20 秒)都去撞一次。
+- **同一个更新只推一次**:推送状态记的是上游 tip。推过之后,上下游没动就不会再推;上游又
+ 有新提交(tip 变了)时会再推一条。
+- **「立即检查」不发通知**:点这个按钮的人正看着结果,没必要再给自己推一条群消息。那次
+ 检查只写结果,没推的那批提交留给下一次定时检查推。
+
+### 部署前提:镜像里要有 git
+
+`python:3.11-slim` 不带 git,`Dockerfile` 里已显式安装。**因此这次更新需要重建镜像**:
+
+```bash
+docker compose build && ./deploy.sh
+```
+
+`./deploy.sh` 只重建前端,不会重建镜像。漏了这步的话,检查会报「未找到 git 命令」——
+界面上看得见,不会静默。
+
+---
+
## 三、必须知道的限制
### 1. 「新增评论」是近似值 —— 最重要的一条
@@ -459,6 +507,9 @@ GET /api/monitor/webhook 通知配置状态(**只返回打码
POST /api/monitor/webhook 保存 Webhook 地址
DELETE /api/monitor/webhook 删除 Webhook
POST /api/monitor/webhook/test 发送测试消息
+
+GET /api/monitor/upstream 最近一次上游检查的缓存结果(没查过返回 {})
+POST /api/monitor/upstream/check 立刻检查一次(等 fetch 跑完才返回,**不发通知**)
```
> `task_id` 用**重复参数**而非逗号拼接(`?task_id=1&task_id=2`);不传表示统计全部任务。
diff --git a/tests/test_upstream.py b/tests/test_upstream.py
new file mode 100644
index 0000000..bc7c821
--- /dev/null
+++ b/tests/test_upstream.py
@@ -0,0 +1,420 @@
+# -*- coding: utf-8 -*-
+# Copyright (c) 2025 relakkes@gmail.com
+#
+# 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
diff --git a/webui/src/components/settings/SystemSettingsDialog.tsx b/webui/src/components/settings/SystemSettingsDialog.tsx
index fa4ac71..66d1a9d 100644
--- a/webui/src/components/settings/SystemSettingsDialog.tsx
+++ b/webui/src/components/settings/SystemSettingsDialog.tsx
@@ -12,6 +12,7 @@ import {
} from '@/components/ui/dialog'
import { WebhookPanel } from '@/components/monitor/WebhookPanel'
import { ChangePassword, SettingField } from '@/components/settings/SettingFields'
+import { UpstreamPanel } from '@/components/settings/UpstreamPanel'
import { useSettings, useUpdateSettings } from '@/hooks/useMonitor'
type Draft = Record
+ 本仓库在上游 MediaCrawler 之上加了一整层,这里定期看看上游有没有新提交。 + 改动立即生效,不需要等下一轮采集。 +
++ {isLoading ? '读取中…' : describe(data)} +
+ + {data?.commits && data.commits.length > 0 && ( +