diff --git a/Dockerfile b/Dockerfile index e8abe0a..382de73 100644 --- a/Dockerfile +++ b/Dockerfile @@ -35,6 +35,10 @@ ENV PYTHONUNBUFFERED=1 \ # # The pip mirror is set for the same reason as the apt one: this host's route to # the public index is slow. +# +# git is for the 上游更新检查 (api/monitor/upstream.py): it fetches the upstream +# repository into the mounted checkout to count how far behind this fork is. +# python:slim does not ship git, and nothing else here pulls it in. ARG APT_MIRROR=mirrors.tuna.tsinghua.edu.cn RUN set -eux; \ for f in /etc/apt/sources.list /etc/apt/sources.list.d/debian.sources; do \ @@ -44,7 +48,7 @@ RUN set -eux; \ done; \ apt-get update; \ apt-get install -y --no-install-recommends \ - build-essential pkg-config default-libmysqlclient-dev \ + build-essential pkg-config default-libmysqlclient-dev git \ libgl1 libglib2.0-0 libsm6 libxext6 libxrender1 libxcb1 libgomp1 \ tzdata nodejs npm; \ rm -rf /var/lib/apt/lists/* diff --git a/UPSTREAM.md b/UPSTREAM.md index fa57e4e..ccc0fd8 100644 --- a/UPSTREAM.md +++ b/UPSTREAM.md @@ -18,6 +18,7 @@ api/auth.py WebUI 登录鉴权 api/monitor/* 监控层整体(含 platforms.py 能力矩阵) api/monitor/db.py MySQL 连接层(可回退 SQLite 供测试用) api/monitor/migrate_from_sqlite.py SQLite → MySQL 一次性迁移脚本 +api/monitor/upstream.py 上游更新检查(定时 fetch 上游并比对) api/routers/{auth,monitor,settings}.py api/schemas/{auth,monitor,settings}.py api/services/interpreter.py 解释器探测(uv / .venv / 当前解释器) @@ -25,10 +26,14 @@ webui/src/components/{monitor,settings,auth}/ 新视图 webui/src/components/layout/{PlatformSwitcher,UnwiredPlatformNotice}.tsx webui/src/{hooks/useMonitor.ts,hooks/usePlatform.ts,store/platformStore.ts,lib/monitorFormat.ts,types/monitor.ts} docs/监控功能使用说明.md -tests/test_{auth,settings,platforms,qrlogin,monitor_*}.py +tests/test_{auth,settings,platforms,qrlogin,monitor_*,upstream}.py Dockerfile / .dockerignore / docker-compose.yml 服务器部署用 ``` +> `api/monitor/upstream.py` 要调 `git`,而 `python:3.11-slim` 不带它 —— Dockerfile 里为此 +> **显式装了 git**。改了 Dockerfile 就必须重建镜像(`docker compose build`),`./deploy.sh` +> 只重建前端,不重建镜像。 + 其中 `api/monitor/qrlogin.py` + `webui/.../QrLoginPanel.tsx` 是**服务器专用**的扫码登录: 那台机器上 Chrome 跑在 Xvfb 里,`show_qrcode` 调的 PIL `Image.show()` 需要桌面看图程序, 服务器没有,二维码会无处可去。所以改成用 CDP 把二维码从页面里读出来交给前端 `` 显示。 @@ -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 @@ -44,6 +45,18 @@ export function SystemSettingsDialog({ [data], ) + // 上游那几项单独成块,其余(时段、CDP)仍归在「调度」下。按前缀分流就够了: + // 注册表是唯一的来源,新增一项上游设置不需要再动这个文件。写成「排除 upstream_」 + // 而不是「只取 active_hours_」,这样以后再加系统设置也不会从界面上凭空消失。 + const upstreamSpecs = useMemo( + () => systemSpecs.filter((spec) => spec.name.startsWith('upstream_')), + [systemSpecs], + ) + const scheduleSpecs = useMemo( + () => systemSpecs.filter((spec) => !spec.name.startsWith('upstream_')), + [systemSpecs], + ) + const dirty = useMemo(() => { if (!data?.values) return {} const changes: Draft = {} @@ -86,7 +99,28 @@ export function SystemSettingsDialog({ 调度器全局只有一套时段规则,因此不按平台区分。

- {systemSpecs.map((spec) => ( + {scheduleSpecs.map((spec) => ( + setDraft((prev) => ({ ...prev, [spec.key]: next }))} + /> + ))} + + +
+
+

+ 上游更新 +

+

+ 本仓库在上游 MediaCrawler 之上加了一整层,这里定期看看上游有没有新提交。 + 改动立即生效,不需要等下一轮采集。 +

+
+ + {upstreamSpecs.map((spec) => ( +
+
{statusBadge(data)}
+ +
+ +

+ {isLoading ? '读取中…' : describe(data)} +

+ + {data?.commits && data.commits.length > 0 && ( +
    + {data.commits.slice(0, 10).map((commit) => ( +
  • + {commit.sha} + + {commit.subject} + +
  • + ))} +
+ )} +
+ ) +} + +function statusBadge(data?: UpstreamStatus) { + if (!data || !data.checked_at) return 尚未检查 + if (!data.ok) return 检查失败 + if ((data.behind ?? 0) > 0) return 落后 {data.behind} 个提交 + return 已是最新 +} + +function describe(data?: UpstreamStatus): string { + if (!data || !data.checked_at) { + return '还没有检查过。开启上面的开关会按间隔自动查,也可以点「立即检查」。' + } + + const when = `上次检查:${formatDateTime(data.checked_at)}(${formatRelative(data.checked_at)})` + if (!data.ok) return `${when};${data.error ?? '未知错误'}` + + const parts = [when, `上游 ${data.branch ?? 'main'} 领先 ${data.behind ?? 0} 个提交`] + // 领先数就是我们自己这一层的规模;合并时要保留的东西,值得一并说出来。 + if (data.ahead) parts.push(`本仓库另有 ${data.ahead} 个自己的提交`) + if (data.notify_error) parts.push(`通知发送失败:${data.notify_error}`) + return parts.join(';') +} diff --git a/webui/src/hooks/useMonitor.ts b/webui/src/hooks/useMonitor.ts index abbd54d..7c449ec 100644 --- a/webui/src/hooks/useMonitor.ts +++ b/webui/src/hooks/useMonitor.ts @@ -352,3 +352,33 @@ export function useTestWebhook() { onError: (error: Error) => toast.error(`发送失败:${error.message}`), }) } + +// --- 上游更新检查 ---------------------------------------------------------- + +export function useUpstreamStatus() { + return useQuery({ + queryKey: ['monitorUpstream'], + queryFn: async () => (await monitorApi.getUpstream()).data, + // 检查本身是每天一次的量级,页面开着时慢点刷就够了。 + refetchInterval: 60_000, + }) +} + +export function useCheckUpstream() { + const queryClient = useQueryClient() + return useMutation({ + mutationFn: () => monitorApi.checkUpstream(), + onSuccess: (response) => { + const data = response.data + if (!data.ok) { + toast.error(`检查失败:${data.error ?? '未知错误'}`) + } else if ((data.behind ?? 0) > 0) { + toast.success(`上游有 ${data.behind} 个新提交`) + } else { + toast.success('已是最新,上游没有新提交') + } + queryClient.invalidateQueries({ queryKey: ['monitorUpstream'] }) + }, + onError: (error: Error) => toast.error(`检查失败:${error.message}`), + }) +} diff --git a/webui/src/lib/api.ts b/webui/src/lib/api.ts index eaf1059..fd8b23b 100644 --- a/webui/src/lib/api.ts +++ b/webui/src/lib/api.ts @@ -16,6 +16,7 @@ import type { ReportResult, SettingsResponse, TaskCreatePayload, + UpstreamStatus, WebhookStatus, } from '@/types/monitor' import type { @@ -278,6 +279,15 @@ export const monitorApi = { setWebhook: (url: string) => api.post('/monitor/webhook', { url }), clearWebhook: () => api.delete('/monitor/webhook'), testWebhook: (url?: string) => api.post('/monitor/webhook/test', { url: url ?? null }), + + /** 最近一次上游检查的缓存结果;没有查过时是空对象。 */ + getUpstream: () => api.get('/monitor/upstream'), + /** + * 立刻检查一次。服务端要等 fetch 跑完才回答,而 axios 默认 30 秒对此不够 —— + * 一条卡住的 git fetch 能拖到两分钟,这里必须单独放长超时,否则会误报失败。 + */ + checkUpstream: () => + api.post('/monitor/upstream/check', null, { timeout: 150_000 }), } /** diff --git a/webui/src/types/monitor.ts b/webui/src/types/monitor.ts index 11a2b19..e7e44b2 100644 --- a/webui/src/types/monitor.ts +++ b/webui/src/types/monitor.ts @@ -379,3 +379,37 @@ export interface SettingsResponse { secrets: Record specs: SettingSpec[] } + +// --- 上游更新检查 ----------------------------------------------------------- + +export interface UpstreamCommit { + sha: string + author: string + /** 提交日期,`YYYY-MM-DD`(git 侧已格式化,不做本地化)。 */ + date: string + subject: string +} + +/** + * 最近一次上游检查的结果,服务端缓存。 + * + * `checked_at` 为空表示还没查过(`GET /monitor/upstream` 返回空对象);`ok` 为 + * false 时 `error` 一定有值 —— 上游不通是常态,那也是一条要显示出来的结论。 + */ +export interface UpstreamStatus { + checked_at?: number | null + remote_url?: string + branch?: string + ok?: boolean + /** 上游有、当前代码没有的提交数 —— 要合的就是这些。 */ + behind?: number + /** 当前代码有、上游没有的提交数 —— 也就是这一层改动自己的规模。 */ + ahead?: number + tip?: string + head?: string + commits?: UpstreamCommit[] + error?: string + /** 本次检查是否推送了通知。 */ + notified?: boolean + notify_error?: string +}