# -*- 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/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 . import adapters 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'> {event.title}') 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 # 链接形状按平台来。抖音的作品是 /video/{id},写死小红书域名的话, # 群里点进去会是一个 404 —— 而这正是通知唯一要它干的事。 url = adapters.adapter(task.platform).note_url(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