Files
MediaCrawler/api/routers/monitor.py
T
butubb b11bbf771a
Deploy VitePress site to Pages / build (push) Canceled after 0s
Deploy VitePress site to Pages / Deploy (push) Canceled after 0s
fix(covers): 封面地址是签名过期而非防盗链,改为本地缓存
实测推翻了之前的诊断。同一批图:

  当天签发的地址 /202610080841/...  -> 200,带不带 Referer 都一样
  隔天的地址     /202610070837/...  -> 403,带不带 Referer 都一样

路径里那段时间戳就是签发时刻。所以这是**过期**,Referer 根本不是那个维度 ——
上一轮加 referrerPolicy 是照着错误结论改的,白改。

修法:
- 采集入库时每轮刷新 cover 地址。原先只在首次入库写一次,旧作品的地址烂在库里,
  而且再怎么重跑也修不回来
- 新增 api/monitor/covers.py:把图下载落盘。图一旦落盘就与签名无关,永远可读
- 下载放在 runner 的 Phase 5(事务已提交之后),不放 ingest —— ingest 的文档写明
  No network,往里塞网络请求会毁掉它可离线测试这一点
- 新增 GET /api/monitor/covers/{note_id} 取图。这条路由带鉴权,封面不会被匿名读走
- service 返回本地地址优先,没有缓存时才退回远程
- 每轮只补一批(60 张):一次跑几百张既慢又会给图床压力,而旧地址本来就在陆续过期,
  分摊到几轮反而更稳

顺带修正 NoteCover 的注释 —— 它写着防盗链,而那个结论已被推翻,留个错的注释比没有更糟。
2026-10-08 08:45:30 +08:00

584 lines
22 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# -*- coding: utf-8 -*-
# Copyright (c) 2025 [email protected]
#
# This file is part of MediaCrawler project.
# Repository: https://github.com/NanmiCoder/MediaCrawler/blob/main/api/routers/monitor.py
# GitHub: https://github.com/NanmiCoder
# Licensed under NON-COMMERCIAL LEARNING LICENSE 1.1
#
# 声明:本代码仅供学习和研究目的使用。使用者应遵守以下原则:
# 1. 不得用于任何商业用途。
# 2. 使用时应遵守目标平台的使用条款和robots.txt规则。
# 3. 不得进行大规模爬取或对平台造成运营干扰。
# 4. 应合理控制请求频率,避免给目标平台带来不必要的负担。
# 5. 不得用于任何非法或不当的用途。
#
# 详细许可条款请参阅项目根目录下的LICENSE文件。
# 使用本代码即表示您同意遵守上述原则和LICENSE中的所有条款。
"""HTTP API for scheduled monitoring tasks."""
from datetime import date, timedelta
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.db import get_session
from ..monitor.platforms import PLATFORM_XHS
from ..monitor.settings import (
cookie_key,
delete_setting,
get_cookie_status,
get_setting,
set_cookie,
set_setting,
)
from ..monitor.models import SETTING_WECOM_WEBHOOK, MonitorTask
from ..schemas.monitor import (
CookiePayload,
MonitorTaskCreate,
MonitorTaskUpdate,
WebhookPayload,
WebhookTestPayload,
)
router = APIRouter(prefix="/monitor", tags=["monitor"])
@router.get("/overview")
async def get_overview(platform: Optional[str] = None):
"""Headline numbers for the dashboard tiles, scoped to one platform."""
async with get_session() as session:
return await service.overview(session, platform)
# ---------------------------------------------------------------------------
# Tasks
# ---------------------------------------------------------------------------
@router.get("/tasks")
async def list_tasks(platform: Optional[str] = None):
async with get_session() as session:
return {"tasks": await service.list_tasks(session, platform)}
@router.post("/tasks", status_code=201)
async def create_task(payload: MonitorTaskCreate):
async with get_session() as session:
try:
task = await service.create_task(session, payload.model_dump())
except service.TargetParseError as exc:
raise HTTPException(status_code=400, detail=str(exc))
except ValueError as exc:
raise HTTPException(status_code=400, detail=str(exc))
return {"id": task.id, "message": "Monitoring task created"}
@router.patch("/tasks/{task_id}")
async def update_task(task_id: int, payload: MonitorTaskUpdate):
async with get_session() as session:
try:
await service.update_task(session, task_id, payload.model_dump(exclude_unset=True))
except service.TargetParseError as exc:
raise HTTPException(status_code=400, detail=str(exc))
except ValueError as exc:
raise HTTPException(status_code=404, detail=str(exc))
return {"message": "Monitoring task updated"}
@router.delete("/tasks/{task_id}")
async def delete_task(task_id: int):
async with get_session() as session:
try:
await service.delete_task(session, task_id)
except ValueError as exc:
raise HTTPException(status_code=404, detail=str(exc))
return {"message": "Monitoring task deleted"}
@router.post("/tasks/{task_id}/run")
async def run_task_now(task_id: int):
"""Queue a run immediately and return; the crawl itself takes minutes."""
async with get_session() as session:
task = await session.get(MonitorTask, task_id)
if task is None:
raise HTTPException(status_code=404, detail=f"Task {task_id} not found")
if not any(target.enabled for target in task.targets):
raise HTTPException(status_code=400, detail="Task has no enabled targets")
service.trigger_manual_run(task_id)
return {"message": "Run queued"}
@router.get("/tasks/{task_id}/runs")
async def list_runs(task_id: int, limit: int = Query(default=50, ge=1, le=500)):
async with get_session() as session:
return {"runs": await service.list_runs(session, task_id, limit=limit)}
# ---------------------------------------------------------------------------
# Collected data
# ---------------------------------------------------------------------------
@router.get("/notes")
async def list_notes(
task_id: Optional[int] = None,
only_new: bool = False,
limit: int = Query(default=200, ge=1, le=2000),
platform: Optional[str] = None,
):
async with get_session() as session:
return {
"notes": await service.list_notes(session, task_id, only_new, limit, platform)
}
@router.get("/notes/{note_id}/series")
async def note_series(note_id: str, task_id: Optional[int] = None):
"""Metric time series for a single note."""
async with get_session() as session:
return {"series": await service.note_series(session, note_id, task_id)}
@router.get("/comments")
async def list_comments(
task_id: Optional[int] = None,
note_id: Optional[str] = None,
group_by: Optional[str] = Query(
default=None, description="传 note 则按作品分组返回,便于阅读"
),
limit: int = Query(default=200, ge=1, le=2000),
platform: Optional[str] = None,
):
"""Comments, each carrying the work it belongs to.
``note_id`` filters to one work; ``group_by=note`` returns them bucketed per
work instead of as a flat stream.
"""
async with get_session() as session:
comments = await service.list_comments(session, task_id, note_id, limit, platform)
if group_by != "note":
return {"comments": comments, "total": len(comments)}
buckets: Dict[str, Dict[str, Any]] = {}
for comment in comments:
bucket = buckets.setdefault(
comment["note_id"],
{
"note_id": comment["note_id"],
"note_title": comment["note_title"],
"note_cover": comment["note_cover"],
"note_url": comment["note_url"],
"comments": [],
},
)
bucket["comments"].append(comment)
ordered = sorted(
buckets.values(),
key=lambda group: group["comments"][0]["first_seen_at"],
reverse=True,
)
return {"groups": ordered, "total": len(comments)}
@router.get("/comment-notes")
async def list_comment_notes(task_id: Optional[int] = None, platform: Optional[str] = None):
"""Works that have comments, newest first, with counts.
Feeds the comment filter dropdown so the operator can pick by title.
"""
async with get_session() as session:
return {"notes": await service.comment_note_groups(session, task_id, platform)}
@router.get("/events")
async def list_events(
task_id: Optional[int] = None,
type: Optional[str] = None,
since_id: Optional[int] = None,
limit: int = Query(default=200, ge=1, le=2000),
platform: Optional[str] = None,
):
async with get_session() as session:
events = await service.list_events(session, task_id, type, since_id, limit, platform)
return {"events": events, "latest_id": events[0]["id"] if events else since_id}
@router.post("/events/read")
async def mark_events_read(task_id: Optional[int] = None):
async with get_session() as session:
count = await service.mark_events_read(session, task_id)
return {"marked": count}
# ---------------------------------------------------------------------------
# Cookie / login health
# ---------------------------------------------------------------------------
@router.get("/cookie")
async def get_cookie_endpoint(platform: str = Query(default=PLATFORM_XHS)):
"""Cookie health only -- deliberately never returns the cookie value.
``platform`` defaults to Xiaohongshu so existing callers keep working; the
key it reads is the namespaced one.
"""
async with get_session() as session:
return await get_cookie_status(session, platform)
@router.post("/cookie")
async def set_cookie_endpoint(payload: CookiePayload, platform: str = Query(default=PLATFORM_XHS)):
async with get_session() as session:
await set_cookie(session, payload.cookie.strip(), platform)
return {"message": "Cookie saved"}
@router.delete("/cookie")
async def clear_cookie_endpoint(platform: str = Query(default=PLATFORM_XHS)):
async with get_session() as session:
await delete_setting(session, cookie_key(platform))
return {"message": "Cookie cleared"}
# ---------------------------------------------------------------------------
# QR login
# ---------------------------------------------------------------------------
# These drive the browser already listening on the CDP debug port, which is the
# same browser -- and therefore the same profile -- that monitor runs attach to.
# Scanning once is what makes later unattended runs logged in.
@router.get("/covers/{note_id}")
async def get_cover(note_id: str):
"""作品封面,从本地缓存读。
**为什么不让前端直连图床**:图床地址是带签名、会过期的 —— 实测隔天即 403,
而且带不带 Referer 都一样,所以那是过期而不是防盗链。本地那份与签名无关。
这个路由是带鉴权的(整条 monitor 路由都挂了 require_auth),所以封面不会被
匿名读走;前端用同源的 <img> 请求会自动带上会话 cookie。
"""
path = covers.find_cached(note_id)
if path is None:
raise HTTPException(status_code=404, detail="封面未缓存")
media_types = {
".jpg": "image/jpeg",
".png": "image/png",
".webp": "image/webp",
".gif": "image/gif",
".heic": "image/heic",
}
return FileResponse(
path,
media_type=media_types.get(path.suffix.lower(), "application/octet-stream"),
# 本地文件不会变(note_id 唯一),让浏览器自己缓存,省掉重复请求。
headers={"Cache-Control": "private, max-age=86400"},
)
@router.post("/login/qr")
async def start_qr_login(platform: str = Query(default=PLATFORM_XHS)):
"""Open the login page in the CDP browser and return its QR code.
A server deployment has no display (Chrome sits under Xvfb), so the code is
surfaced here for the operator to scan instead of in a desktop window that
does not exist.
"""
try:
return await qrlogin.start(platform)
except ValueError as exc:
raise HTTPException(status_code=400, detail=str(exc))
except RuntimeError as exc:
raise HTTPException(status_code=502, detail=str(exc))
@router.get("/login/qr")
async def get_qr_login():
"""Poll the live session: waiting -> success / expired / error."""
return await qrlogin.status()
@router.delete("/login/qr")
async def cancel_qr_login():
"""Drop our tab and stop polling."""
return await qrlogin.cancel()
@router.get("/login/state")
async def get_login_state(force: bool = Query(default=False)):
"""Ask the browser itself whether it is signed in.
Deliberately separate from the QR session above. That session is in-memory and
dies with the process -- a redeploy is enough -- so "am I logged in?" must not
hinge on it, or a successful scan looks like nothing happened.
``force`` reloads the page first, for when the login may have lapsed somewhere
else and the page's copy of the state is stale.
"""
return await qrlogin.check_login_state(force=force)
# ---------------------------------------------------------------------------
# Report
# ---------------------------------------------------------------------------
async def _resolve_scope(
session, task_ids: Optional[List[int]], platform: Optional[str]
) -> Optional[List[int]]:
"""Combine an explicit task selection with an optional platform filter.
``None`` means "no restriction"; an explicit list is intersected with the
platform's tasks so a stale selection cannot leak another platform's data
into a scoped report.
"""
if platform is None:
return task_ids
platform_ids = set(await service.platform_task_ids(session, platform))
if task_ids is None:
return list(platform_ids)
return [task for task in task_ids if task in platform_ids]
@router.get("/export")
async def export_data(
kind: str = Query(..., description="notes | comments | report"),
task_id: Optional[List[int]] = Query(default=None),
note_id: Optional[str] = None,
start_date: Optional[str] = Query(default=None, description="YYYY-MM-DD,report 用"),
end_date: Optional[str] = Query(default=None, description="YYYY-MM-DD,report 用"),
days: int = Query(default=7, ge=1, le=365),
file_format: str = Query(default="csv", alias="format", description="csv | xlsx"),
platform: Optional[str] = None,
):
"""Download collected data as CSV or Excel.
Reached by the browser as a navigation (``window.open``), which cannot carry
an Authorization header -- this is one of the reasons the session lives in a
cookie.
"""
if kind not in ("notes", "comments", "report"):
raise HTTPException(status_code=400, detail="kind 必须是 notes / comments / report")
if file_format not in ("csv", "xlsx"):
raise HTTPException(status_code=400, detail="format 必须是 csv 或 xlsx")
async with get_session() as session:
scoped = await _resolve_scope(session, task_id, platform)
if kind == "notes":
single_task = scoped[0] if scoped and len(scoped) == 1 else None
rows = await service.list_notes(session, single_task, False, 5000, platform)
elif kind == "comments":
single_task = scoped[0] if scoped and len(scoped) == 1 else None
rows = await service.list_comments(session, single_task, note_id, 5000, platform)
else:
try:
end_day = date.fromisoformat(end_date) if end_date else date.today()
start_day = (
date.fromisoformat(start_date)
if start_date
else end_day - timedelta(days=days - 1)
)
except ValueError:
raise HTTPException(status_code=400, detail="日期格式应为 YYYY-MM-DD")
# Not `report = ...`: that would make `report` a local name for the
# whole function and shadow the module import on this very line.
report_data = await report.build_report(session, scoped, start_day, end_day)
rows = report_data["rows"]
if not rows:
raise HTTPException(status_code=404, detail="该范围内没有数据可导出")
columns = _export_columns(kind)
stamp = date.today().isoformat()
# ASCII on purpose: a non-ASCII filename needs RFC 5987 encoding in
# Content-Disposition, and the plain `filename="..."` form used below would
# mangle it.
filename = f"export_{kind}_{stamp}.{file_format}"
if file_format == "xlsx":
payload = _to_xlsx(rows, columns, kind)
media_type = "application/vnd.openxmlformats-officedocument.spreadsheetml.sheet"
else:
payload = _to_csv(rows, columns)
media_type = "text/csv; charset=utf-8"
return Response(
content=payload,
media_type=media_type,
headers={"Content-Disposition": f'attachment; filename="{filename}"'},
)
def _export_columns(kind: str) -> List[tuple[str, str]]:
"""(key, header) pairs per export kind."""
if kind == "notes":
return [
("note_id", "作品ID"),
("title", "标题"),
("note_url", "链接"),
("liked_count", "点赞"),
("comment_count", "评论"),
("collected_count", "收藏"),
("share_count", "分享"),
("liked_count_delta", "点赞增量"),
("comment_count_delta", "评论增量"),
("first_seen_at", "首次发现"),
("last_seen_at", "最近采集"),
]
if kind == "comments":
return [
("note_title", "所属作品"),
("note_id", "作品ID"),
("comment_id", "评论ID"),
("content", "内容"),
("nickname", "昵称"),
("like_count", "点赞"),
("sub_comment_count", "子评论数"),
("create_time", "发布时间"),
("first_seen_at", "首次发现"),
]
return [
("date", "日期"),
("new_notes", "新增作品"),
("new_comments", "新增评论"),
("liked_count_delta", "点赞增量"),
("comment_count_delta", "评论增量"),
("collected_count_delta", "收藏增量"),
("share_count_delta", "分享增量"),
]
def _cell(value: Any) -> Any:
if value is None:
return ""
if isinstance(value, (list, dict)):
return ", ".join(str(v) for v in value) if isinstance(value, list) else str(value)
return value
def _to_csv(rows: List[Dict[str, Any]], columns: List[tuple[str, str]]) -> bytes:
import csv
import io
buffer = io.StringIO()
writer = csv.writer(buffer)
writer.writerow([header for _, header in columns])
for row in rows:
writer.writerow([_cell(row.get(key)) for key, _ in columns])
# utf-8-sig: without the BOM Excel opens Chinese CSV as mojibake, which is
# the single most common complaint about CSV exports here.
return buffer.getvalue().encode("utf-8-sig")
def _to_xlsx(rows: List[Dict[str, Any]], columns: List[tuple[str, str]], sheet: str) -> bytes:
import io
from openpyxl import Workbook
workbook = Workbook()
worksheet = workbook.active
worksheet.title = {"notes": "作品", "comments": "评论"}.get(sheet, "报表")
worksheet.append([header for _, header in columns])
for row in rows:
worksheet.append([_cell(row.get(key)) for key, _ in columns])
output = io.BytesIO()
workbook.save(output)
return output.getvalue()
@router.get("/report")
async def get_report(
task_id: Optional[List[int]] = Query(
default=None, description="Repeat to include several tasks; omit for all"
),
start_date: Optional[str] = Query(default=None, description="YYYY-MM-DD"),
end_date: Optional[str] = Query(default=None, description="YYYY-MM-DD"),
days: int = Query(default=7, ge=1, le=365, description="Window used when dates are omitted"),
platform: Optional[str] = None,
):
"""Daily new-content counts and interaction deltas for the selected tasks."""
try:
end_day = date.fromisoformat(end_date) if end_date else date.today()
start_day = date.fromisoformat(start_date) if start_date else end_day - timedelta(days=days - 1)
except ValueError:
raise HTTPException(status_code=400, detail="日期格式应为 YYYY-MM-DD")
if start_day > end_day:
raise HTTPException(status_code=400, detail="开始日期不能晚于结束日期")
async with get_session() as session:
scoped = await _resolve_scope(session, task_id, platform)
return await report.build_report(session, scoped, start_day, end_day)
# ---------------------------------------------------------------------------
# WeCom webhook
# ---------------------------------------------------------------------------
def _mask_webhook(url: str) -> str:
"""Show enough of the URL to recognise it, without exposing the robot key."""
if not url:
return ""
key_marker = "key="
index = url.find(key_marker)
if index == -1:
return url[:12] + "..." if len(url) > 12 else url
prefix = url[: index + len(key_marker)]
key = url[index + len(key_marker) :]
if len(key) <= 8:
return prefix + "*" * len(key)
return f"{prefix}{key[:4]}...{key[-4:]}"
@router.get("/webhook")
async def get_webhook():
async with get_session() as session:
url = (await get_setting(session, SETTING_WECOM_WEBHOOK)) or ""
return {"configured": bool(url), "masked": _mask_webhook(url)}
@router.post("/webhook")
async def set_webhook(payload: WebhookPayload):
url = payload.url.strip()
if url and "qyapi.weixin.qq.com" not in url:
# Catches the common mistake of pasting a group-chat invite or the app
# URL instead of the robot webhook.
raise HTTPException(
status_code=400,
detail="这不像企业微信机器人 Webhook 地址(应包含 qyapi.weixin.qq.com)",
)
async with get_session() as session:
await set_setting(session, SETTING_WECOM_WEBHOOK, url)
return {"message": "Webhook 已保存" if url else "Webhook 已清空", "configured": bool(url)}
@router.delete("/webhook")
async def clear_webhook():
async with get_session() as session:
await delete_setting(session, SETTING_WECOM_WEBHOOK)
return {"message": "Webhook 已删除"}
@router.post("/webhook/test")
async def test_webhook(payload: WebhookTestPayload):
"""Send a test message so the user can verify the robot works before relying on it."""
async with get_session() as session:
url = payload.url.strip() if payload.url else await notify.get_webhook_url(session)
ok, detail = await notify.send_wecom(
url, "**综合采集平台 通知测试**\n> 如果你看到这条消息,说明 Webhook 配置成功。"
)
if not ok:
raise HTTPException(status_code=400, detail=detail)
return {"message": detail}