Files
butubb f6ddc46d62
Deploy VitePress site to Pages / build (push) Canceled after 0s
Deploy VitePress site to Pages / Deploy (push) Canceled after 0s
feat(notify): 通知拆成「新作品」与「异常」两个开关,异常默认开
问题:cookie 过期导致任务失败,但没有任何通知。查下来不是代码问题 ——
notify_enabled 在两个任务上都是 False,而它默认就是关的,事件(run_failed)
也确实生成了,只卡在最后一道闸门。

但那个默认值是错的。代码里的理由是「一条任务列表都推到一个群会很快变吵,所以默认静默」,
这个理由对新作品成立(可能每轮都有),对失败不成立:一次登录态失效意味着这个任务事实上
已经死了,而你不会知道,直到某天发现数据停在几周前。最该被告知的就是这种情况。

现在拆开:
- notify_enabled  —— 推送新作品,可能每轮都有,默认关
- notify_failures —— 推送异常(登录失效/运行失败/没抓到数据),默认开

事件按开关过滤(build_run_message):只勾了「新作品」的任务不该因为一次失败被推消息,
反之亦然,否则拆开开关就没有意义。已有任务由 _ensure_columns 补上 notify_failures=1,
所以会自动开始收到异常推送。

列名 notify_enabled 是历史遗留(它早先是唯一的通知开关),语义已收窄为「新作品」,
用注释写明,不做列重命名 —— 那需要单独的迁移,不值为一个内部工具做。
2026-10-09 13:37:02 +08:00

215 lines
7.3 KiB
Python
Raw Permalink 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/monitor/notify.py
# GitHub: https://github.com/NanmiCoder
# Licensed under NON-COMMERCIAL LEARNING LICENSE 1.1
#
# 声明:本代码仅供学习和研究目的使用。使用者应遵守以下原则:
# 1. 不得用于任何商业用途。
# 2. 使用时应遵守目标平台的使用条款和robots.txt规则。
# 3. 不得进行大规模爬取或对平台造成运营干扰。
# 4. 应合理控制请求频率,避免给目标平台带来不必要的负担。
# 5. 不得用于任何非法或不当的用途。
#
# 详细许可条款请参阅项目根目录下的LICENSE文件。
# 使用本代码即表示您同意遵守上述原则和LICENSE中的所有条款。
"""Push notifications via a WeCom (企业微信) group robot webhook.
Two rules shape this module:
* **One message per run, not per event.** A run that finds twenty new notes must
produce one summary, not twenty pushes.
* **A failed push never fails the crawl.** Notification is best-effort: the run's
data is already committed by the time we get here, so every error is logged
and swallowed.
"""
import json
from typing import Optional
import httpx
from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncSession
from tools.time_util import get_current_timestamp
from .models import (
EVENT_AUTH_FAILURE,
EVENT_NEW_NOTE,
EVENT_NO_DATA,
EVENT_RUN_FAILED,
MonitorEvent,
MonitorRun,
MonitorTask,
)
from .settings import get_setting
# Short on purpose: the scheduler awaits the run, so a hanging webhook would
# stall every other task behind it.
WEBHOOK_TIMEOUT_SECONDS = 10.0
# Only these event types are worth interrupting someone for. NO_DATA is included
# because a run that fetched nothing at all is always anomalous -- a creator
# always has *some* notes -- even when the login is not the culprit.
NOTIFIABLE_EVENT_TYPES = (
EVENT_AUTH_FAILURE,
EVENT_RUN_FAILED,
EVENT_NO_DATA,
EVENT_NEW_NOTE,
)
# WeCom markdown is a limited subset; coloured text is the one bit of flair it
# supports and it makes failures stand out in a busy group chat.
_COLOR_WARNING = "warning"
_COLOR_INFO = "info"
async def send_wecom(webhook_url: str, content: str) -> tuple[bool, str]:
"""Post a markdown message to a WeCom group robot.
Returns (ok, detail) rather than raising, so callers can surface the reason
in the UI when the user clicks "send test".
"""
if not webhook_url:
return False, "Webhook 未配置"
payload = {"msgtype": "markdown", "markdown": {"content": content}}
try:
async with httpx.AsyncClient(timeout=WEBHOOK_TIMEOUT_SECONDS) as client:
response = await client.post(webhook_url, json=payload)
response.raise_for_status()
body = response.json()
except httpx.HTTPError as exc:
return False, f"请求失败:{exc}"
except json.JSONDecodeError:
return False, "返回内容不是合法 JSON,请检查 Webhook 地址"
# WeCom answers 200 with a non-zero errcode on failure.
errcode = body.get("errcode")
if errcode != 0:
return False, f"企业微信返回 errcode={errcode} {body.get('errmsg', '')}"
return True, "发送成功"
async def get_webhook_url(session: AsyncSession) -> str:
from .models import SETTING_WECOM_WEBHOOK
return (await get_setting(session, SETTING_WECOM_WEBHOOK)) or ""
async def build_run_message(
session: AsyncSession,
task: MonitorTask,
run: MonitorRun,
) -> Optional[str]:
"""Compose one markdown summary for a finished run, or None if nothing to say.
事件按开关过滤:只勾了「新作品」的任务,不该因为一次失败被推消息,反之亦然 ——
否则拆开这两个开关就没有意义了。
"""
allowed = []
if task.notify_enabled:
allowed.append(EVENT_NEW_NOTE)
if task.notify_failures:
allowed.extend([EVENT_AUTH_FAILURE, EVENT_RUN_FAILED, EVENT_NO_DATA])
if not allowed:
return None
events = list(
(
await session.scalars(
select(MonitorEvent)
.where(
MonitorEvent.run_id == run.id,
MonitorEvent.type.in_(allowed),
)
.order_by(MonitorEvent.id)
)
).all()
)
if not events:
return None
failures = [
e for e in events if e.type in (EVENT_AUTH_FAILURE, EVENT_RUN_FAILED, EVENT_NO_DATA)
]
new_notes = [e for e in events if e.type == EVENT_NEW_NOTE]
lines: list[str] = []
if failures:
# Word the header from what actually happened, not from whether new notes
# accompanied it: a login outage usually brings no new notes either.
unavailable = any(e.type == EVENT_NO_DATA for e in failures) and not any(
e.type in (EVENT_AUTH_FAILURE, EVENT_RUN_FAILED) for e in failures
)
header = "监控任务未抓到数据" if unavailable else "监控任务异常"
lines.append(f"**⚠️ {header}:{task.name}**")
for event in failures:
lines.append(f'> <font color="{_COLOR_WARNING}">{event.title}</font>')
else:
lines.append(f"**📢 监控任务有新作品:{task.name}**")
if new_notes:
lines.append(f"> 新增作品 **{len(new_notes)}** 篇")
# Cap the listing: a first-ever run or a long gap can produce a lot, and
# a wall of text is worse than a count.
for event in new_notes[:10]:
payload = _load_payload(event.payload_json)
title = payload.get("title") or event.target_id
note_id = payload.get("note_id") or event.target_id
url = f"https://www.xiaohongshu.com/explore/{note_id}"
lines.append(f"> [{title}]({url})")
if len(new_notes) > 10:
lines.append(f"> …等共 {len(new_notes)} 篇")
if run.is_baseline:
lines.append("> (首轮基线,未计入新增统计)")
return "\n".join(lines)
def _load_payload(raw: str) -> dict:
try:
payload = json.loads(raw or "{}")
except json.JSONDecodeError:
return {}
return payload if isinstance(payload, dict) else {}
async def notify_run(session: AsyncSession, task: MonitorTask, run: MonitorRun) -> Optional[str]:
"""Push a summary for a finished run if the task opted in.
Returns the message that was sent, or None. Never raises.
"""
try:
# 两个开关是分开的:只开「异常」不该因为新作品而发消息,反之亦然。
if not (task.notify_enabled or task.notify_failures):
return None
webhook_url = await get_webhook_url(session)
if not webhook_url:
return None
message = await build_run_message(session, task, run)
if not message:
return None
ok, detail = await send_wecom(webhook_url, message)
if not ok:
print(f"[monitor.notify] task {task.id} push failed: {detail}")
return None
task.last_notified_at = get_current_timestamp()
return message
except Exception as exc: # pragma: no cover - notification must never break a run
print(f"[monitor.notify] unexpected error: {exc}")
return None