# -*- 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_ingest.py # GitHub: https://github.com/NanmiCoder # Licensed under NON-COMMERCIAL LEARNING LICENSE 1.1 # # 声明:本代码仅供学习和研究目的使用。使用者应遵守以下原则: # 1. 不得用于任何商业用途。 # 2. 使用时应遵守目标平台的使用条款和robots.txt规则。 # 3. 不得进行大规模爬取或对平台造成运营干扰。 # 4. 应合理控制请求频率,避免给目标平台带来不必要的负担。 # 5. 不得用于任何非法或不当的用途。 # # 详细许可条款请参阅项目根目录下的LICENSE文件。 # 使用本代码即表示您同意遵守上述原则和LICENSE中的所有条款。 """Offline tests for the monitoring ingest/diff layer. These run without network, browser or login and cover the correctness caveats that matter most: baseline suppression, count parsing, NULL-vs-zero, the posted/seen comment split, idempotency, and the silent-cookie-failure signal. """ import json from datetime import date from pathlib import Path from typing import Any, Dict, List, Optional import pytest import pytest_asyncio from sqlalchemy import select from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker, create_async_engine from sqlalchemy.pool import StaticPool from tools.time_util import get_current_timestamp from api.monitor import adapters from api.monitor.ingest import describe_exit_code, ingest_run, parse_count from api.monitor.models import ( EVENT_AUTH_FAILURE, EVENT_METRIC_DELTA, EVENT_NEW_COMMENT_POSTED, EVENT_NEW_COMMENT_SEEN, EVENT_NEW_NOTE, EVENT_NO_DATA, EVENT_RUN_FAILED, MODE_CREATOR, MonitorBase, MonitorComment, MonitorEvent, MonitorNote, MonitorNoteMetric, MonitorRun, MonitorTask, RUN_FAILED, RUN_PARTIAL, RUN_SUCCESS, ) @pytest_asyncio.fixture async def db(): """An isolated in-memory monitoring database.""" 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 db_session: yield db_session await engine.dispose() async def _make_task(db: AsyncSession, **overrides) -> MonitorTask: defaults = dict( name="test task", 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, created_at=0, updated_at=0, ) defaults.update(overrides) task = MonitorTask(**defaults) db.add(task) await db.flush() return task async def _make_run( db: AsyncSession, task: MonitorTask, started_at: int, exit_code: Optional[int] = 0, ) -> MonitorRun: run = MonitorRun( task_id=task.id, trigger="manual", status=RUN_SUCCESS, phase=task.mode, save_data_path="", queued_at=started_at, not_before=0, started_at=started_at, exit_code=exit_code, ) db.add(run) await db.flush() return run def _write_run_dir( root: Path, notes: List[Dict[str, Any]], comments: Optional[List[Dict[str, Any]]] = None, subdir: str = "xhs", ) -> Path: """Write a run's jsonl output in the crawler's own layout. ``subdir`` 是**爬虫**落盘的目录名,不是监控层的平台 id —— 抖音那边这两者不同 (平台 id 是 ``dy``、目录是 ``douyin``),所以必须能分开指定,否则测不出那个差异。 """ jsonl_dir = root / subdir / "jsonl" jsonl_dir.mkdir(parents=True, exist_ok=True) contents = jsonl_dir / "creator_contents_2026-01-01.jsonl" contents.write_text( "\n".join(json.dumps(n, ensure_ascii=False) for n in notes), encoding="utf-8", ) if comments is not None: comment_file = jsonl_dir / "creator_comments_2026-01-01.jsonl" comment_file.write_text( "\n".join(json.dumps(c, ensure_ascii=False) for c in comments), encoding="utf-8", ) return root def _note(note_id: str, liked: Any = "10", **extra) -> Dict[str, Any]: record = { "note_id": note_id, "title": f"title-{note_id}", "note_url": f"https://www.xiaohongshu.com/explore/{note_id}", "image_list": "https://img/cover.jpg", "creator_hash": "hash", "time": 1700000000000, "liked_count": liked, "comment_count": "1", "collected_count": "1", "share_count": "1", } record.update(extra) return record def _comment(comment_id: str, note_id: str, create_time: int, **extra) -> Dict[str, Any]: record = { "comment_id": comment_id, "note_id": note_id, "content": f"content-{comment_id}", "nickname": "u***r", "creator_hash": "hash", "create_time": create_time, "like_count": "0", "sub_comment_count": 0, "parent_comment_id": "", } record.update(extra) return record def _dy_note(aweme_id: str, liked: Any = "10", **extra) -> Dict[str, Any]: """抖音作品记录 —— 键名照抄 store/douyin/__init__.py 的落盘字段。 重点在于**没有** ``note_id``:抖音叫 ``aweme_id``。这一条差异没映射好,就是 每条记录都被 ingest 悄悄 continue 掉、一条不剩。 """ record = { "aweme_id": aweme_id, "aweme_type": "0", "title": f"title-{aweme_id}", "desc": f"title-{aweme_id}", # 抖音给的是**秒**(实测 1790574515 = 2026-09-28),小红书给毫秒。落库统一 # 换算成毫秒,这个 fixture 必须照真实形态写,否则测不出单位问题。 "create_time": 1790574515, "creator_hash": "hash", "nickname": "u***r", "liked_count": liked, "collected_count": "1", "comment_count": "1", "share_count": "1", "aweme_url": f"https://www.douyin.com/video/{aweme_id}", "cover_url": "https://img/cover.jpg", } record.update(extra) return record def _dy_comment( comment_id: str, aweme_id: str, create_time: int, **extra ) -> Dict[str, Any]: record = { "comment_id": comment_id, "create_time": create_time, "aweme_id": aweme_id, "content": f"content-{comment_id}", "creator_hash": "hash", "nickname": "u***r", "sub_comment_count": "0", "like_count": "0", # 抖音顶层评论的父 id 是字符串 "0",不是空串。 "parent_comment_id": "0", } record.update(extra) return record async def _events(db: AsyncSession, event_type: Optional[str] = None) -> List[MonitorEvent]: stmt = select(MonitorEvent) if event_type: stmt = stmt.where(MonitorEvent.type == event_type) return list((await db.scalars(stmt)).all()) # -------------------------------------------------------------------------- # parse_count # -------------------------------------------------------------------------- class TestParseCount: @pytest.mark.parametrize( "raw,expected", [ ("1234", 1234), ("1.2万", 12000), ("1.2w", 12000), ("3亿", 300000000), ("1,234", 1234), (42, 42), ], ) def test_parses_platform_formats(self, raw, expected): assert parse_count(raw) == expected @pytest.mark.parametrize("raw", ["", None, "暂无", "-", "abc", True]) def test_unparseable_values_return_none(self, raw): assert parse_count(raw) is None # -------------------------------------------------------------------------- # Exit codes # -------------------------------------------------------------------------- class TestExitCodeStorage: """Guards a bug that only showed up when the data moved to MySQL. Windows reports process failures as unsigned 32-bit NTSTATUS values (0xC0000142 = 3221225794). That overflows MySQL's signed INT, while SQLite's dynamic typing accepted it happily -- so the column silently worked until a real migration hit it with real data. """ def test_column_is_bigint_not_int(self): from sqlalchemy import BigInteger from api.monitor.models import MonitorRun column_type = MonitorRun.__table__.c.exit_code.type assert isinstance(column_type, BigInteger), ( f"exit_code must be BigInteger to hold unsigned 32-bit codes, got {column_type!r}" ) @pytest.mark.asyncio async def test_an_ntstatus_value_round_trips(self, db): task = await _make_task(db) run = await _make_run(db, task, started_at=1000, exit_code=3221225794) await db.commit() stored = await db.scalar( select(MonitorRun.exit_code).where(MonitorRun.id == run.id) ) assert stored == 3221225794 class TestDescribeExitCode: def test_windows_status_code_is_decoded(self): """3221225794 is 0xC0000142, which is meaningless without decoding.""" message = describe_exit_code(3221225794) assert "0xC0000142" in message assert "DLL_INIT_FAILED" in message def test_negative_signed_form_is_also_decoded(self): # Python may hand back the signed form depending on how it was launched. assert "0xC0000142" in describe_exit_code(-1073741502) def test_unknown_code_degrades_to_the_raw_number(self): assert describe_exit_code(1) == "Crawler exited with code 1" # -------------------------------------------------------------------------- # Notes # -------------------------------------------------------------------------- class TestNoteIngest: @pytest.mark.asyncio async def test_baseline_run_emits_no_new_note_events(self, db, tmp_path): task = await _make_task(db) run = await _make_run(db, task, started_at=1000) _write_run_dir(tmp_path, [_note("n1"), _note("n2")], comments=[]) result = await ingest_run(db, run, task, tmp_path) assert result.status == RUN_SUCCESS assert result.is_baseline is True assert result.new_notes == 2 # Everything is "new" on the first run; emitting that would be pure noise. assert await _events(db, EVENT_NEW_NOTE) == [] assert len(list((await db.scalars(select(MonitorNote))).all())) == 2 @pytest.mark.asyncio async def test_an_empty_run_does_not_establish_a_baseline(self, db, tmp_path): """A run that fetched nothing observed nothing, so it is not a baseline. Otherwise the first crawl that actually works reports every work as newly discovered. """ task = await _make_task(db) empty_run = await _make_run(db, task, started_at=1000) (tmp_path / "empty").mkdir(parents=True, exist_ok=True) await ingest_run(db, empty_run, task, tmp_path / "empty") real_run = await _make_run(db, task, started_at=2000) result = await ingest_run( db, real_run, task, _write_run_dir(tmp_path / "ok", [_note("n1")], comments=[]) ) assert result.is_baseline is True assert await _events(db, EVENT_NEW_NOTE) == [] @pytest.mark.asyncio async def test_second_run_reports_only_the_added_note(self, db, tmp_path): task = await _make_task(db) first_dir = _write_run_dir(tmp_path / "run1", [_note("n1")], comments=[]) run1 = await _make_run(db, task, started_at=1000) await ingest_run(db, run1, task, first_dir) second_dir = _write_run_dir(tmp_path / "run2", [_note("n1"), _note("n2")], comments=[]) run2 = await _make_run(db, task, started_at=2000) result = await ingest_run(db, run2, task, second_dir) assert result.is_baseline is False assert result.new_notes == 1 events = await _events(db, EVENT_NEW_NOTE) assert len(events) == 1 assert events[0].target_id == "n2" assert events[0].run_id == run2.id class TestMetricSnapshots: @pytest.mark.asyncio async def test_delta_event_emitted_when_like_count_changes(self, db, tmp_path): task = await _make_task(db) run1 = await _make_run(db, task, started_at=1000) await ingest_run(db, run1, task, _write_run_dir(tmp_path / "r1", [_note("n1", "100")], comments=[])) run2 = await _make_run(db, task, started_at=2000) await ingest_run(db, run2, task, _write_run_dir(tmp_path / "r2", [_note("n1", "150")], comments=[])) events = await _events(db, EVENT_METRIC_DELTA) assert len(events) == 1 payload = json.loads(events[0].payload_json) assert payload["deltas"]["liked_count"] == {"from": 100, "to": 150, "delta": 50} @pytest.mark.asyncio async def test_no_delta_when_nothing_changed(self, db, tmp_path): task = await _make_task(db) run1 = await _make_run(db, task, started_at=1000) await ingest_run(db, run1, task, _write_run_dir(tmp_path / "r1", [_note("n1", "100")], comments=[])) run2 = await _make_run(db, task, started_at=2000) await ingest_run(db, run2, task, _write_run_dir(tmp_path / "r2", [_note("n1", "100")], comments=[])) assert await _events(db, EVENT_METRIC_DELTA) == [] @pytest.mark.asyncio async def test_unparseable_count_is_null_not_zero(self, db, tmp_path): task = await _make_task(db) run = await _make_run(db, task, started_at=1000) await ingest_run(db, run, task, _write_run_dir(tmp_path, [_note("n1", "暂无")], comments=[])) metric = await db.scalar(select(MonitorNoteMetric).where(MonitorNoteMetric.note_id == "n1")) # Zero would forge a large negative delta on the next comparison. assert metric.liked_count is None assert metric.raw_liked_count == "暂无" @pytest.mark.asyncio async def test_no_delta_when_previous_value_was_unparseable(self, db, tmp_path): task = await _make_task(db) run1 = await _make_run(db, task, started_at=1000) await ingest_run(db, run1, task, _write_run_dir(tmp_path / "r1", [_note("n1", "暂无")], comments=[])) run2 = await _make_run(db, task, started_at=2000) await ingest_run(db, run2, task, _write_run_dir(tmp_path / "r2", [_note("n1", "50")], comments=[])) assert await _events(db, EVENT_METRIC_DELTA) == [] @pytest.mark.asyncio async def test_metric_snapshot_survives_across_runs(self, db, tmp_path): """The crawler's own DB store overwrites metrics; ours must not.""" task = await _make_task(db) for index, liked in enumerate(["100", "150", "300"]): run = await _make_run(db, task, started_at=1000 * (index + 1)) await ingest_run( db, run, task, _write_run_dir(tmp_path / f"r{index}", [_note("n1", liked)], comments=[]) ) snapshots = list( ( await db.scalars( select(MonitorNoteMetric) .where(MonitorNoteMetric.note_id == "n1") .order_by(MonitorNoteMetric.run_id) ) ).all() ) assert [s.liked_count for s in snapshots] == [100, 150, 300] # -------------------------------------------------------------------------- # Comments # -------------------------------------------------------------------------- class TestCommentIngest: @pytest.mark.asyncio async def test_posted_vs_seen_split_by_create_time(self, db, tmp_path): task = await _make_task(db) # Baseline establishes the seen-set; no events on the first run. run1 = await _make_run(db, task, started_at=1000) await ingest_run( db, run1, task, _write_run_dir(tmp_path / "r1", [_note("n1")], comments=[_comment("c1", "n1", create_time=500)]), ) assert await _events(db, EVENT_NEW_COMMENT_POSTED) == [] # c2 was published after run1 started -> genuinely new. # c3 is old but only just surfaced in the top-N window -> seen, not posted. run2 = await _make_run(db, task, started_at=2000) await ingest_run( db, run2, task, _write_run_dir( tmp_path / "r2", [_note("n1")], comments=[ _comment("c1", "n1", create_time=500), _comment("c2", "n1", create_time=2500), _comment("c3", "n1", create_time=100), ], ), ) posted = await _events(db, EVENT_NEW_COMMENT_POSTED) seen = await _events(db, EVENT_NEW_COMMENT_SEEN) assert len(posted) == 1 assert json.loads(posted[0].payload_json)["comment_id"] == "c2" assert len(seen) == 1 assert json.loads(seen[0].payload_json)["comment_id"] == "c3" @pytest.mark.asyncio async def test_comments_not_ingested_when_disabled(self, db, tmp_path): task = await _make_task(db, enable_comments=False) run = await _make_run(db, task, started_at=1000) result = await ingest_run( db, run, task, _write_run_dir(tmp_path, [_note("n1")], comments=[_comment("c1", "n1", 500)]), ) assert result.new_comments == 0 # -------------------------------------------------------------------------- # Failure handling # -------------------------------------------------------------------------- class TestFailureHandling: @pytest.mark.asyncio async def test_nonzero_exit_is_a_failure(self, db, tmp_path): task = await _make_task(db) run = await _make_run(db, task, started_at=1000, exit_code=1) _write_run_dir(tmp_path, [_note("n1")], comments=[]) result = await ingest_run(db, run, task, tmp_path) assert result.status == RUN_FAILED assert len(await _events(db, EVENT_RUN_FAILED)) == 1 # A crashed run must not touch the seen-set. assert await db.scalar(select(MonitorNote.id)) is None @pytest.mark.asyncio async def test_zero_notes_with_exit_zero_is_a_suspected_auth_failure(self, db, tmp_path): """The silent-cookie-failure signature: exit 0 but nothing fetched. A real bad-cookie run writes no output file at all, which is why the exit code has to be checked before the files are. """ task = await _make_task(db) run = await _make_run(db, task, started_at=1000, exit_code=0) tmp_path.mkdir(parents=True, exist_ok=True) result = await ingest_run(db, run, task, tmp_path) assert result.status == RUN_PARTIAL assert len(await _events(db, EVENT_AUTH_FAILURE)) == 1 assert await _events(db, EVENT_RUN_FAILED) == [] @pytest.mark.asyncio async def test_no_data_is_not_blamed_on_the_cookie_when_a_sibling_succeeded( self, db, tmp_path ): """A task that just worked proves the login is fine; do not cry wolf.""" healthy = await _make_task(db, name="healthy") healthy_run = await _make_run(db, healthy, started_at=get_current_timestamp()) await ingest_run( db, healthy_run, healthy, _write_run_dir(tmp_path / "ok", [_note("n1")], comments=[]), ) task = await _make_task(db, name="suspect") run = await _make_run(db, task, started_at=get_current_timestamp()) (tmp_path / "empty").mkdir(parents=True, exist_ok=True) result = await ingest_run(db, run, task, tmp_path / "empty") assert result.status == RUN_PARTIAL assert await _events(db, EVENT_NO_DATA) != [] assert await _events(db, EVENT_AUTH_FAILURE) == [] @pytest.mark.asyncio async def test_empty_contents_file_is_also_an_auth_failure(self, db, tmp_path): task = await _make_task(db) run = await _make_run(db, task, started_at=1000, exit_code=0) _write_run_dir(tmp_path, [], comments=[]) result = await ingest_run(db, run, task, tmp_path) assert result.status == RUN_PARTIAL assert len(await _events(db, EVENT_AUTH_FAILURE)) == 1 # -------------------------------------------------------------------------- # Idempotency # -------------------------------------------------------------------------- class TestIdempotency: @pytest.mark.asyncio async def test_reingesting_the_same_data_adds_nothing(self, db, tmp_path): task = await _make_task(db) run_dir = _write_run_dir( tmp_path, [_note("n1"), _note("n2")], comments=[_comment("c1", "n1", 500)] ) run1 = await _make_run(db, task, started_at=1000) await ingest_run(db, run1, task, run_dir) notes_after_first = len(list((await db.scalars(select(MonitorNote))).all())) # A retry of the same crawl content must not duplicate rows or events. run2 = await _make_run(db, task, started_at=2000) result = await ingest_run(db, run2, task, run_dir) assert result.new_notes == 0 assert result.new_comments == 0 assert len(list((await db.scalars(select(MonitorNote))).all())) == notes_after_first class TestDouyinIngest: """抖音的产物形状与小红书不同 —— 这里钉住「不会被静默丢掉」。 这一组存在的理由,是这个改动最危险的失败模式:字段名或目录名没对上时,ingest 不报错,只是**一条都不入库**,然后被当成「疑似登录失效」报出去。 """ async def _ingest( self, db, tmp_path, notes, comments=None, platform="dy", subdir="douyin", ): task = await _make_task(db, platform=platform) run = await _make_run(db, task, started_at=1) _write_run_dir(tmp_path, notes, comments, subdir=subdir) result = await ingest_run(db, run, task, tmp_path) return task, run, result @pytest.mark.asyncio async def test_notes_are_ingested_under_their_douyin_field_names(self, db, tmp_path): aweme_id = "7525082444551310602" _task, _run, result = await self._ingest(db, tmp_path, [_dy_note(aweme_id)]) note = await db.scalar(select(MonitorNote)) assert note is not None, "抖音作品被静默丢弃了 —— 多半是 aweme_id 没映射到 note_id" assert note.note_id == aweme_id assert note.note_url == f"https://www.douyin.com/video/{aweme_id}" assert note.cover == "https://img/cover.jpg" assert note.source_kind == "0" # 秒 → 毫秒,换算过才对。 assert note.published_at == 1790574515 * 1000 assert result.notes_fetched == 1 @pytest.mark.asyncio async def test_timestamps_are_normalised_to_milliseconds(self, db, tmp_path): """抖音的时间戳是**秒**,小红书是毫秒 —— 差 1000 倍,必须换算。 不换算的话,2026 年的作品会显示成 1970 年。这是实测踩到的:抖音作品的 「发布日期」列显示成 1970-01-22(1790574515 被当成毫秒就是 21 天后)。 """ aweme_id = "7525082444551310602" await self._ingest(db, tmp_path, [_dy_note(aweme_id)]) note = await db.scalar(select(MonitorNote)) assert note.published_at == 1790574515 * 1000 # 落库的是毫秒,展示层才不用关心来源;乘完应该在 2026 年,不是 1970。 assert date.fromtimestamp(note.published_at / 1000).year == 2026 @pytest.mark.asyncio async def test_the_artifact_directory_is_not_the_platform_id(self, db, tmp_path): """目录名与平台 id 不一致,是这套适配里最反直觉的一条。 抖音的平台 id 是 ``dy``,而爬虫把产物写在 ``douyin/`` 下。把它钉在这里, 是为了让「顺手改成一致」这件事会在测试里红掉,而不是让 ingest 悄悄读 0 条。 """ assert adapters.artifact_dir("dy") == "douyin" @pytest.mark.asyncio async def test_writing_into_the_platform_id_directory_reads_nothing(self, db, tmp_path): """反面:产物落在 ``dy/`` 下时一条都读不到 —— 这正是映射要解决的问题。""" _task, _run, result = await self._ingest(db, tmp_path, [_dy_note("1")], subdir="dy") assert result.notes_fetched == 0 @pytest.mark.asyncio async def test_misplaced_output_is_blamed_on_the_directory_not_the_login( self, db, tmp_path ): """产物其实抓到了,只是目录名不对 —— 不该报成「疑似登录失效」。 这是最难查的一类故障:登录是好的、数据也抓到了,但报出来的现象和登录失效 一模一样,会把人指去查完全错误的方向。 """ _task, run, _result = await self._ingest( db, tmp_path, [_dy_note("1")], subdir="dy" ) assert any("目录" in event.title for event in await _events(db, EVENT_NO_DATA)) assert await _events(db, EVENT_AUTH_FAILURE) == [] assert run.error_message and "dy" in run.error_message @pytest.mark.asyncio async def test_comments_are_linked_through_aweme_id(self, db, tmp_path): aweme_id = "7525082444551310602" _task, _run, result = await self._ingest( db, tmp_path, [_dy_note(aweme_id)], comments=[_dy_comment("c1", aweme_id, 500)], ) comment = await db.scalar(select(MonitorComment)) assert comment is not None, "抖音评论被静默丢弃了 —— 多半是 aweme_id 没映射" assert comment.note_id == aweme_id assert result.comments_fetched == 1 @pytest.mark.asyncio async def test_a_top_level_parent_of_zero_becomes_empty(self, db, tmp_path): """抖音顶层评论的父 id 是 "0";原样存进去,前端会多出一堆悬空的父节点。""" aweme_id = "7525082444551310602" await self._ingest( db, tmp_path, [_dy_note(aweme_id)], comments=[ _dy_comment("c1", aweme_id, 500), _dy_comment("c2", aweme_id, 600, parent_comment_id="c1"), ], ) by_id = {c.comment_id: c for c in (await db.scalars(select(MonitorComment))).all()} assert by_id["c1"].parent_comment_id == "" assert by_id["c2"].parent_comment_id == "c1" @pytest.mark.asyncio async def test_the_four_metrics_need_no_mapping(self, db, tmp_path): """四个指标键两边同名 —— 抖音作品照样进 monitor_note_metric,差分照常。""" aweme_id = "7525082444551310602" task = await _make_task(db, platform="dy") _write_run_dir(tmp_path, [_dy_note(aweme_id, liked="100")], subdir="douyin") run1 = await _make_run(db, task, started_at=1) await ingest_run(db, run1, task, tmp_path) metric = await db.scalar(select(MonitorNoteMetric)) assert metric is not None and metric.liked_count == 100 _write_run_dir(tmp_path, [_dy_note(aweme_id, liked="150")], subdir="douyin") run2 = await _make_run(db, task, started_at=2) await ingest_run(db, run2, task, tmp_path) assert len(await _events(db, EVENT_METRIC_DELTA)) == 1 class TestNicknameRefresh: """已入库的评论,昵称要跟着重新采集的值走。 评论是去重后直接跳过的,若不刷新,脱敏开关一改(或评论者改了昵称),老数据就永远 停在旧值上 —— 而重采是唯一能拿到新值的途径。 """ @pytest.mark.asyncio async def test_an_existing_comment_gets_its_nickname_refreshed(self, db, tmp_path): task = await _make_task(db) _write_run_dir(tmp_path, [_note("n1")], comments=[_comment("c1", "n1", 500)]) run1 = await _make_run(db, task, started_at=1) await ingest_run(db, run1, task, tmp_path) assert (await db.scalar(select(MonitorComment))).nickname == "u***r" _write_run_dir( tmp_path, [_note("n1")], comments=[_comment("c1", "n1", 500, nickname="未脱敏的新昵称")], ) run2 = await _make_run(db, task, started_at=2) await ingest_run(db, run2, task, tmp_path) comment = await db.scalar(select(MonitorComment)) assert comment.nickname == "未脱敏的新昵称" # 去重的语义没变:同一条评论不该被插成两行。 assert ( len(list((await db.scalars(select(MonitorComment))).all())) == 1 )