# -*- 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_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中的所有条款。 """Tests for the WeCom notification layer. The webhook is stubbed, so nothing here touches the network. """ import json import pytest import pytest_asyncio from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker, create_async_engine from sqlalchemy.pool import StaticPool from api.monitor import notify from api.monitor.models import ( EVENT_AUTH_FAILURE, EVENT_METRIC_DELTA, EVENT_NEW_NOTE, EVENT_NEW_COMMENT_POSTED, MODE_CREATOR, SETTING_WECOM_WEBHOOK, MonitorBase, MonitorEvent, MonitorRun, MonitorTask, RUN_SUCCESS, ) from api.monitor.settings import set_setting WEBHOOK = "https://qyapi.weixin.qq.com/cgi-bin/webhook/send?key=abc123" @pytest_asyncio.fixture async def db(): engine = create_async_engine("sqlite+aiosqlite://", poolclass=StaticPool) async with engine.begin() as conn: await conn.run_sync(MonitorBase.metadata.create_all) factory = async_sessionmaker(engine, class_=AsyncSession, expire_on_commit=False) async with factory() as session: yield session await engine.dispose() async def _seed(db: AsyncSession, notify_enabled: bool = True): task = MonitorTask( name="竞品监控", platform="xhs", mode=MODE_CREATOR, enabled=True, interval_minutes=60, max_notes_count=20, enable_comments=True, max_comments_count=50, run_timeout_seconds=3600, notify_enabled=notify_enabled, created_at=0, updated_at=0, ) db.add(task) await db.flush() run = MonitorRun( task_id=task.id, trigger="scheduled", status=RUN_SUCCESS, phase=MODE_CREATOR, save_data_path="", queued_at=0, not_before=0, max_comments_count=50, ) db.add(run) await db.flush() return task, run def _add_event(db, task, run, event_type, title, payload=None, severity="info"): db.add( MonitorEvent( task_id=task.id, run_id=run.id, type=event_type, severity=severity, target_kind="note", target_id="note-1", title=title, payload_json=json.dumps(payload or {}, ensure_ascii=False), created_at=0, ) ) # -------------------------------------------------------------------------- # Message building # -------------------------------------------------------------------------- class TestBuildRunMessage: @pytest.mark.asyncio async def test_no_notifiable_events_means_no_message(self, db): task, run = await _seed(db) # Metric deltas are not something anyone wants pushed. _add_event(db, task, run, EVENT_METRIC_DELTA, "点赞 10→20") _add_event(db, task, run, EVENT_NEW_COMMENT_POSTED, "新评论") await db.flush() assert await notify.build_run_message(db, task, run) is None @pytest.mark.asyncio async def test_new_notes_are_listed_with_links(self, db): task, run = await _seed(db) _add_event( db, task, run, EVENT_NEW_NOTE, "新作品:标题A", payload={"note_id": "abc123", "title": "标题A"}, ) await db.flush() message = await notify.build_run_message(db, task, run) assert "竞品监控" in message assert "新增作品 **1** 篇" in message assert "标题A" in message assert "https://www.xiaohongshu.com/explore/abc123" in message @pytest.mark.asyncio async def test_long_note_lists_are_truncated(self, db): """A first run can find dozens; a wall of text is worse than a count.""" task, run = await _seed(db) for index in range(14): _add_event( db, task, run, EVENT_NEW_NOTE, f"新作品:{index}", payload={"note_id": f"n{index}", "title": f"标题{index}"}, ) await db.flush() message = await notify.build_run_message(db, task, run) assert "新增作品 **14** 篇" in message assert "标题0" in message assert "标题13" not in message assert "等共 14 篇" in message @pytest.mark.asyncio async def test_failure_is_reported_as_a_warning(self, db): task, run = await _seed(db) _add_event( db, task, run, EVENT_AUTH_FAILURE, "疑似登录态失效:本次未抓到任何作品", severity="error", ) await db.flush() message = await notify.build_run_message(db, task, run) assert "异常" in message assert "登录态失效" in message assert notify._COLOR_WARNING in message @pytest.mark.asyncio async def test_baseline_runs_say_so(self, db): task, run = await _seed(db) run.is_baseline = True _add_event(db, task, run, EVENT_NEW_NOTE, "新作品", payload={"note_id": "x", "title": "t"}) await db.flush() message = await notify.build_run_message(db, task, run) assert "基线" in message # -------------------------------------------------------------------------- # notify_run gating # -------------------------------------------------------------------------- class TestNotifyRunGating: @pytest.mark.asyncio async def test_disabled_task_is_skipped(self, db, monkeypatch): task, run = await _seed(db, notify_enabled=False) _add_event(db, task, run, EVENT_NEW_NOTE, "新作品", payload={"note_id": "x", "title": "t"}) await set_setting(db, SETTING_WECOM_WEBHOOK, WEBHOOK) await db.flush() called = [] monkeypatch.setattr(notify, "send_wecom", lambda *a, **k: called.append(a) or _ok()) assert await notify.notify_run(db, task, run) is None assert called == [] @pytest.mark.asyncio async def test_missing_webhook_is_skipped(self, db, monkeypatch): task, run = await _seed(db, notify_enabled=True) _add_event(db, task, run, EVENT_NEW_NOTE, "新作品", payload={"note_id": "x", "title": "t"}) await db.flush() called = [] monkeypatch.setattr(notify, "send_wecom", lambda *a, **k: called.append(a) or _ok()) assert await notify.notify_run(db, task, run) is None assert called == [] @pytest.mark.asyncio async def test_successful_push_records_the_timestamp(self, db, monkeypatch): task, run = await _seed(db, notify_enabled=True) _add_event(db, task, run, EVENT_NEW_NOTE, "新作品", payload={"note_id": "x", "title": "t"}) await set_setting(db, SETTING_WECOM_WEBHOOK, WEBHOOK) await db.flush() monkeypatch.setattr(notify, "send_wecom", lambda *a, **k: _ok()) message = await notify.notify_run(db, task, run) assert message is not None # Lets the UI answer "why did I not get a push for this run?". assert task.last_notified_at is not None @pytest.mark.asyncio async def test_push_failure_never_raises(self, db, monkeypatch): """A broken webhook must not take down the crawl that just succeeded.""" task, run = await _seed(db, notify_enabled=True) _add_event(db, task, run, EVENT_NEW_NOTE, "新作品", payload={"note_id": "x", "title": "t"}) await set_setting(db, SETTING_WECOM_WEBHOOK, WEBHOOK) await db.flush() async def _boom(*args, **kwargs): raise RuntimeError("network exploded") monkeypatch.setattr(notify, "send_wecom", _boom) assert await notify.notify_run(db, task, run) is None async def _ok(): return True, "发送成功" # -------------------------------------------------------------------------- # send_wecom # -------------------------------------------------------------------------- class _FakeResponse: def __init__(self, payload): self._payload = payload def raise_for_status(self): return None def json(self): return self._payload class _FakeClient: """Captures the request and replays a canned WeCom reply.""" last_payload = None def __init__(self, reply=None, error=None): self._reply = reply if reply is not None else {"errcode": 0, "errmsg": "ok"} self._error = error def __call__(self, *args, **kwargs): return self async def __aenter__(self): return self async def __aexit__(self, *exc): return False async def post(self, url, json=None): if self._error: raise self._error type(self).last_payload = json return _FakeResponse(self._reply) class TestSendWecom: @pytest.mark.asyncio async def test_missing_url_is_reported(self): ok, detail = await notify.send_wecom("", "hi") assert ok is False assert "未配置" in detail @pytest.mark.asyncio async def test_success(self, monkeypatch): monkeypatch.setattr(notify.httpx, "AsyncClient", _FakeClient()) ok, detail = await notify.send_wecom(WEBHOOK, "**标题**\n> 内容") assert ok is True assert detail == "发送成功" # WeCom expects a markdown message envelope. assert _FakeClient.last_payload["msgtype"] == "markdown" assert _FakeClient.last_payload["markdown"]["content"] == "**标题**\n> 内容" @pytest.mark.asyncio async def test_nonzero_errcode_is_a_failure(self, monkeypatch): """WeCom answers HTTP 200 even when it rejects the message.""" monkeypatch.setattr( notify.httpx, "AsyncClient", _FakeClient(reply={"errcode": 93000, "errmsg": "invalid webhook url"}), ) ok, detail = await notify.send_wecom(WEBHOOK, "hi") assert ok is False assert "93000" in detail @pytest.mark.asyncio async def test_network_error_is_returned_not_raised(self, monkeypatch): import httpx monkeypatch.setattr( notify.httpx, "AsyncClient", _FakeClient(error=httpx.ConnectError("boom")), ) ok, detail = await notify.send_wecom(WEBHOOK, "hi") assert ok is False assert "请求失败" in detail