diff --git a/api/monitor/douyin_api.py b/api/monitor/douyin_api.py index 9fc6f2f..3fc98a1 100644 --- a/api/monitor/douyin_api.py +++ b/api/monitor/douyin_api.py @@ -425,6 +425,12 @@ async def author_profile(sec_user_id: str, *, cookie: str = "") -> Dict[str, Any f"接口没返回用户数据(status_code={payload.get('status_code')})" ) return { + # 自报家门。快照表的唯一键是 (任务, creator_hash, 轮次),而作品是靠 + # `anonymize_user_id(author.uid)` 得到这个哈希的 —— 这里走同一条路,两边才对得上, + # 否则快照会和作品分成两个人,界面上永远查不到。 + "creator_hash": anonymize_user_id( + str(user.get("uid") or user.get("sec_uid") or "") + ), "nickname": user.get("nickname") or "", "unique_id": user.get("unique_id") or "", "fans": _as_int(user.get("follower_count")), diff --git a/api/monitor/douyin_fetch.py b/api/monitor/douyin_fetch.py index 6b49d3b..5344678 100644 --- a/api/monitor/douyin_fetch.py +++ b/api/monitor/douyin_fetch.py @@ -32,7 +32,7 @@ import json from datetime import datetime from pathlib import Path -from typing import Any, Dict, Iterable, List, Sequence +from typing import Any, Dict, Iterable, List, Optional, Sequence from tools import utils @@ -65,6 +65,9 @@ async def collect( """ notes: List[Dict[str, Any]] = [] comments: List[Dict[str, Any]] = [] + # 博主**账号级**指标(粉丝 / 总获赞 / 作品数)。作品列表之外单独要一次, + # 只有博主模式才有 —— 作品模式的目标是一件作品,没有"这个博主是谁"可问。 + profiles: List[Dict[str, Any]] = [] errors: List[str] = [] # **整个 collect 只去重一次的、跨目标的集合**:退化路径会把「库里已知的全部作品」 @@ -89,6 +92,9 @@ async def collect( videos = await _creator_works( external_id, limit, known_aweme_ids, seen_aweme, cookie, errors ) + profile = await _creator_profile(external_id, videos, cookie, errors) + if profile is not None: + profiles.append(profile) for video in videos: aweme_id = video.get("aweme_id") @@ -107,7 +113,7 @@ async def collect( except douyin_api.DouyinApiError as exc: errors.append(f"拉取作品 {aweme_id} 的评论失败:{exc}") - jsonl_dir = _write_artifacts(out_dir, platform, mode, notes, comments) + jsonl_dir = _write_artifacts(out_dir, platform, mode, notes, comments, profiles) return { "notes": len(notes), "comments": len(comments), @@ -116,6 +122,35 @@ async def collect( } +async def _creator_profile( + sec_user_id: str, + videos: Sequence[Dict[str, Any]], + cookie: str, + errors: List[str], +) -> Optional[Dict[str, Any]]: + """问一次博主的账号级指标。拿不到就算了 —— **不能因为顺手的附加信息失败, + 就把这一轮本来采到的作品也判成失败。** + + creator_hash 优先取作品自带的那个:作品是靠 ``anonymize_user_id(author.uid)`` 得到 + 哈希的,而快照表和作品必须对得上号,否则界面上永远查不出这个博主的粉丝数。只有当一件 + 作品都没采到时(列表被挡且没有已知作品可刷新),才退回资料接口自己算的哈希 —— + 那种情况下也只剩它了。 + """ + try: + profile = await douyin_api.author_profile(sec_user_id, cookie=cookie) + except douyin_api.DouyinApiError as exc: + errors.append(f"拉取博主 {sec_user_id} 的资料失败:{exc}") + return None + + if videos: + profile["creator_hash"] = videos[0].get("creator_hash") or profile["creator_hash"] + if not profile.get("creator_hash"): + # 哈希都算不出来的快照没人能查到,落下去只是垃圾。 + errors.append(f"博主 {sec_user_id} 的资料里没有可用的身份标识,跳过账号指标") + return None + return profile + + async def _creator_works( sec_user_id: str, limit: int, @@ -151,6 +186,7 @@ def _write_artifacts( mode: str, notes: List[Dict[str, Any]], comments: List[Dict[str, Any]], + profiles: Sequence[Dict[str, Any]] = (), ) -> Path: """按爬虫那套目录与文件名写 jsonl。 @@ -167,6 +203,10 @@ def _write_artifacts( # 评论文件即使没有评论也建出来:ingest 靠「文件在不在」区分「这一轮没评论」和 # 「这一轮什么都没抓到」,两种情况的含义完全不同。 _write_jsonl(jsonl_dir / f"{kind}_comments_{date}.jsonl", comments) + # 博主资料同样无条件写:空文件表示"问了但没问到",没有文件表示"这次根本没问" + # (作品模式)。两者在 ingest 那边走的是同一条路(都不落快照),但留空文件能让 + # 事后翻 run 目录时看出到底问没问过。 + _write_jsonl(jsonl_dir / f"{kind}_profile_{date}.jsonl", list(profiles)) return jsonl_dir diff --git a/api/monitor/ingest.py b/api/monitor/ingest.py index b6a12f2..ba682d1 100644 --- a/api/monitor/ingest.py +++ b/api/monitor/ingest.py @@ -56,6 +56,7 @@ from .models import ( EVENT_NO_DATA, EVENT_RUN_FAILED, MonitorComment, + MonitorCreatorStat, MonitorEvent, MonitorNote, MonitorNoteMetric, @@ -219,6 +220,19 @@ def find_run_files( ) +def find_profile_files(out_dir: Path, platform: str = PLATFORM_XHS) -> List[Path]: + """博主**账号级**指标那几行 jsonl(``creator_profile_*.jsonl``)。 + + 单独一个函数而不是往 ``find_run_files`` 的返回值里塞第三个列表:那个返回值的两个 + 位置是有意义的(contents/comments),加一个会把所有调用点和解包语句都牵动一遍, + 而这份产物是**可选**的 —— 小红书那条路(爬虫进程)根本不产生它。 + """ + jsonl_dir = out_dir / adapters.artifact_dir(platform) / "jsonl" + if not jsonl_dir.is_dir(): + return [] + return sorted(jsonl_dir.glob("*_profile_*.jsonl")) + + def _misplaced_output_dirs(out_dir: Path, expected: str) -> List[str]: """在 out_dir 下找「有产物、但目录名不是期望的那个」的目录。 @@ -593,6 +607,49 @@ async def _ingest_comments( return new_count +async def _ingest_creator_stats( + session: AsyncSession, + run: MonitorRun, + records: Sequence[Dict[str, Any]], +) -> int: + """把这一轮问到的博主账号级指标落成快照,返回条数。 + + 和作品指标一样是**每轮一条**:账号级的粉丝数是缓慢变化的量,「今天比昨天多了 300」 + 才是有用的信号,单看一个绝对值没有意义 —— 所以这里只管记,分析交给查询端。 + + **没解析出来的值留 NULL,不写 0**:0 在趋势图上是一条砸到底的线,和「不知道」完全是 + 两回事(见 ``parse_count`` 的注释)。 + """ + now = get_current_timestamp() + written = 0 + seen: set = set() + + for record in records: + creator_hash = str(record.get("creator_hash") or "").strip() + if not creator_hash or creator_hash in seen: + # 一个任务可以配多个目标,退化路径下它们可能指向同一个博主 —— 而唯一键是 + # (task_id, creator_hash, run_id),重复插入会撞键把整轮炸掉。 + continue + seen.add(creator_hash) + + session.add( + MonitorCreatorStat( + task_id=run.task_id, + run_id=run.id, + creator_hash=creator_hash, + nickname=str(record.get("nickname") or "")[:128], + fans=parse_count(record.get("fans")), + total_favorited=parse_count(record.get("total_favorited")), + works_count=parse_count(record.get("works")), + following=parse_count(record.get("following")), + captured_at=now, + ) + ) + written += 1 + + return written + + async def ingest_run( session: AsyncSession, run: MonitorRun, @@ -639,6 +696,16 @@ async def ingest_run( run.notes_fetched = len(contents) run.comments_fetched = len(comments) + # 账号级快照**在「一条作品都没采到」的早退之前**落。博主的粉丝数并不会因为他这个 + # 月的新作品列表被风控挡住就不存在 —— 那正是最该看到「粉丝还在涨、但新作品没在发现」 + # 的时刻,跳过它等于在最需要它的那轮把数据丢掉。 + profiles = [ + record + for path in find_profile_files(out_dir, task.platform) + for record in _read_jsonl(path) + ] + await _ingest_creator_stats(session, run, profiles) + # A bad cookie does NOT fail the process: XHS cookie login is never validated, # so an unauthenticated session just returns zero notes with exit 0 -- and # usually does not even create an output file. Treating that as "the creator diff --git a/api/monitor/models.py b/api/monitor/models.py index 17370a0..20c5bb4 100644 --- a/api/monitor/models.py +++ b/api/monitor/models.py @@ -343,6 +343,64 @@ class MonitorCreatorAlias(MonitorBase): updated_at: Mapped[int] = mapped_column(BigInteger, nullable=False) +class MonitorCreatorStat(MonitorBase): + """博主的**账号级**快照:粉丝数 / 总获赞 / 作品数 / 关注数。 + + 这是作品列表给不了的东西:作品级指标说"这一条视频涨了多少赞",账号级说"这个人 + 整个账号的粉丝是在涨还是在掉"。两者不互相替代。 + + 粒度取 ``(任务, 博主, 轮次)``,和作品指标一样的形状 —— 于是趋势、差分、报表那套 + 现成的逻辑换个表就能用。 + + 目前**只有抖音**会写它:小红书那条走的是爬虫子进程,而它的 ``save_creator()`` 在 + 教学版里是空函数,根本没落过创作者资料。所以表里只有抖音的博主。 + """ + + __tablename__ = "monitor_creator_stat" + __table_args__ = ( + UniqueConstraint( + "task_id", "creator_hash", "run_id", name="uq_creator_stat" + ), + ) + + id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True) + task_id: Mapped[int] = mapped_column( + ForeignKey("monitor_task.id", ondelete="CASCADE"), nullable=False, index=True + ) + run_id: Mapped[int] = mapped_column(Integer, nullable=False, index=True) + creator_hash: Mapped[str] = mapped_column(String(64), nullable=False, index=True) + nickname: Mapped[str] = mapped_column(String(128), nullable=False, default="") + + # 都可能为 None:平台没给就留空,**不要伪造成 0** —— 0 是"掉到零",和"不知道" + # 在趋势图上是完全不同的两回事。 + fans: Mapped[Optional[int]] = mapped_column(BigInteger) + total_favorited: Mapped[Optional[int]] = mapped_column(BigInteger) + works_count: Mapped[Optional[int]] = mapped_column(BigInteger) + following: Mapped[Optional[int]] = mapped_column(BigInteger) + + captured_at: Mapped[int] = mapped_column(BigInteger, nullable=False, index=True) + + +class MonitorNoteAlias(MonitorBase): + """给**作品**起的备注。 + + 和 ``MonitorCreatorAlias`` 是一对:博主那条回答"这是谁",这条回答"这条我要盯着"。 + 键取 ``(platform, note_id)`` —— 作品 id 本身就带平台语义,但显式带上 platform 才能和 + 博主备注用同一套查询形状。 + """ + + __tablename__ = "monitor_note_alias" + __table_args__ = ( + UniqueConstraint("platform", "note_id", name="uq_note_alias"), + ) + + id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True) + platform: Mapped[str] = mapped_column(String(16), nullable=False, index=True) + note_id: Mapped[str] = mapped_column(String(128), nullable=False, index=True) + alias: Mapped[str] = mapped_column(String(128), nullable=False, default="") + updated_at: Mapped[int] = mapped_column(BigInteger, nullable=False) + + class MonitorSetting(MonitorBase): """Key/value store. Holds the XHS cookie for unattended runs.""" diff --git a/api/monitor/service.py b/api/monitor/service.py index 658b35f..1ef65bd 100644 --- a/api/monitor/service.py +++ b/api/monitor/service.py @@ -19,7 +19,7 @@ """Task CRUD and dashboard queries for the monitoring layer.""" import asyncio -from typing import Any, Dict, List, Optional +from typing import Any, Dict, List, Optional, Sequence from urllib.parse import parse_qs, urlparse from sqlalchemy import delete, func, select @@ -35,8 +35,10 @@ from .models import ( MODE_NOTE, MonitorComment, MonitorCreatorAlias, + MonitorCreatorStat, MonitorEvent, MonitorNote, + MonitorNoteAlias, MonitorNoteMetric, MonitorRun, MonitorTarget, @@ -389,6 +391,74 @@ def _delta(current: Optional[int], previous: Optional[int]) -> Optional[int]: return current - previous +async def _note_alias_map(session: AsyncSession) -> Dict[tuple, str]: + """``(platform, note_id) -> 作品备注``。和博主备注一个道理,整体读一次。""" + rows = (await session.scalars(select(MonitorNoteAlias))).all() + return {(row.platform, row.note_id): row.alias for row in rows if row.alias} + + +async def set_note_alias( + session: AsyncSession, platform: str, note_id: str, alias: str +) -> None: + """给作品起备注;空串就是删掉这条备注。""" + alias = (alias or "").strip()[:128] + existing = await session.scalar( + select(MonitorNoteAlias).where( + MonitorNoteAlias.platform == platform, + MonitorNoteAlias.note_id == note_id, + ) + ) + + if not alias: + if existing is not None: + await session.delete(existing) + return + + if existing is None: + session.add( + MonitorNoteAlias( + platform=platform, + note_id=note_id, + alias=alias, + updated_at=get_current_timestamp(), + ) + ) + return + + existing.alias = alias + existing.updated_at = get_current_timestamp() + + +async def _latest_creator_stats( + session: AsyncSession, + task_ids: Sequence[int], + creator_hashes: Sequence[str], +) -> Dict[tuple, "MonitorCreatorStat"]: + """``(task_id, creator_hash) -> 最近一条``账号级快照。 + + 按 run_id 而不是 captured_at 取「最近」:和作品指标用的是同一个口径,两者放一起 + 看才不会出现「作品数据来自第 8 轮、粉丝数来自第 9 轮」这种对不上的情况。 + """ + if not task_ids or not creator_hashes: + return {} + + rows = ( + await session.scalars( + select(MonitorCreatorStat) + .where( + MonitorCreatorStat.task_id.in_(list(task_ids)), + MonitorCreatorStat.creator_hash.in_(list(creator_hashes)), + ) + .order_by(MonitorCreatorStat.run_id.desc()) + ) + ).all() + + latest: Dict[tuple, MonitorCreatorStat] = {} + for row in rows: + latest.setdefault((row.task_id, row.creator_hash), row) # 已按 run_id 倒序 + return latest + + async def list_notes( session: AsyncSession, task_id: Optional[int] = None, @@ -442,6 +512,16 @@ async def list_notes( ).all() } aliases = await _creator_alias_map(session) + note_aliases = await _note_alias_map(session) + # 账号级指标(粉丝 / 总获赞 / 作品数)。**挂在作品上一起返回**,因为界面上就是按博主 + # 归组显示的 —— 让前端为了一个组头再发一轮请求没道理。同一个博主的所有作品拿到的是 + # 同一条(键里带 task_id,所以跨任务不会串)。没有的(小红书那条路不产生它)就是 null, + # 前端据此整块不显示,而不是显示一个 0。 + creator_stats = await _latest_creator_stats( + session, + [note.task_id for note in notes], + [note.creator_hash for note in notes if note.creator_hash], + ) for note in notes: series = by_note.get(note.note_id, []) @@ -454,6 +534,8 @@ async def list_notes( if note.first_seen_run_id != latest_run_ids[note.task_id]: continue + stat = creator_stats.get((note.task_id, note.creator_hash)) + result.append( { "task_id": note.task_id, @@ -474,6 +556,17 @@ async def list_notes( "creator_alias": aliases.get( (task_platform.get(note.task_id, ""), note.creator_hash), "" ), + # 这条作品自己的备注。和博主备注是两件事:博主备注回答"这是谁",它回答 + # "这条我要盯着"。 + "note_alias": note_aliases.get( + (task_platform.get(note.task_id, ""), note.note_id), "" + ), + # 博主账号级指标 —— 作品列表给不了的东西。三个值都可能为 null(平台没采 + # 到、或者这条作品来自不产生它的数据源),前端据此整块不画。 + "creator_fans": stat.fans if stat else None, + "creator_total_favorited": stat.total_favorited if stat else None, + "creator_works": stat.works_count if stat else None, + "creator_stats_at": stat.captured_at if stat else None, # 优先给本地缓存地址:远程地址带签名、会过期(实测隔天即 403), # 本地那份不会。没有缓存时才退回远程,至少让图先显示出来。 "cover": covers.cover_url(note.note_id, note.cover), diff --git a/api/routers/monitor.py b/api/routers/monitor.py index da6c828..7ccafb0 100644 --- a/api/routers/monitor.py +++ b/api/routers/monitor.py @@ -39,6 +39,7 @@ from ..monitor.models import SETTING_WECOM_WEBHOOK, MonitorTask from ..schemas.monitor import ( CookiePayload, CreatorAliasPayload, + NoteAliasPayload, MonitorTaskCreate, MonitorTaskUpdate, WebhookPayload, @@ -275,6 +276,22 @@ async def set_creator_alias_endpoint( return {"creator_hash": creator_hash, "alias": payload.alias.strip()} +@router.put("/notes/{note_id}") +async def set_note_alias_endpoint( + note_id: str, + payload: NoteAliasPayload, + platform: str = Query(default=PLATFORM_XHS), +): + """给**作品**起个备注。 + + 和上面那条博主备注是一对:博主备注回答「这个账号是谁」,这条回答「这条作品我要盯着」。 + 两者不能合并 —— 一个博主底下常常只有一两件值得盯的作品。 + """ + async with get_session() as session: + await service.set_note_alias(session, platform, note_id, payload.alias) + return {"note_id": note_id, "alias": payload.alias.strip()} + + # --------------------------------------------------------------------------- # QR login # --------------------------------------------------------------------------- diff --git a/api/schemas/monitor.py b/api/schemas/monitor.py index 94eb3a9..622bdac 100644 --- a/api/schemas/monitor.py +++ b/api/schemas/monitor.py @@ -114,6 +114,12 @@ class CreatorAliasPayload(BaseModel): alias: str = Field(default="", max_length=128) +class NoteAliasPayload(BaseModel): + """给作品起的备注。空串表示清掉这条备注。""" + + alias: str = Field(default="", max_length=128) + + class WebhookPayload(BaseModel): url: str = Field(default="", description="企业微信机器人 Webhook 地址,留空表示停用") diff --git a/tests/test_douyin_fetch.py b/tests/test_douyin_fetch.py index 172cd52..8774871 100644 --- a/tests/test_douyin_fetch.py +++ b/tests/test_douyin_fetch.py @@ -13,6 +13,35 @@ import pytest from api.monitor import douyin_api, douyin_fetch +def _profile(**overrides) -> dict: + profile = { + "creator_hash": "hash", + "nickname": "博主", + "unique_id": "abc", + "fans": 12000, + "total_favorited": 83000, + "works": 42, + "following": 7, + } + profile.update(overrides) + return profile + + +@pytest.fixture(autouse=True) +def fake_profile(monkeypatch): + """每个用例都挡住「问博主资料」这一跳。 + + 它是附加信息,不在任何一条编排路径上,但真发出去就会去连 9222 那个浏览器 —— + 于是所有 creator 用例都会多出一次连接失败、并把 ``errors`` 弄脏。想验它自己的 + 用例再单独覆盖这个 fixture。 + """ + + async def _profile_call(sec_user_id, *, cookie=""): + return _profile() + + monkeypatch.setattr(douyin_api, "author_profile", _profile_call) + + def _video(aweme_id: str, likes: str = "1") -> dict: return { "aweme_id": aweme_id, @@ -272,3 +301,111 @@ class TestDegradation: # 评论拿不到是小事,作品不能跟着丢。 assert result["notes"] == 1 assert any("评论失败" in error for error in result["errors"]) + + +class TestCreatorProfile: + """博主的**账号级**指标 —— 粉丝 / 总获赞 / 作品数。 + + 作品列表给不了这个东西:它说的是一件作品涨了多少赞,不是这个人整个账号的粉丝 + 在涨还是在掉。单独问一次资料接口。 + """ + + @staticmethod + def _profiles(tmp_path): + files = list((tmp_path / "douyin" / "jsonl").glob("*_profile_*.jsonl")) + assert len(files) == 1 + return _read(files[0]) + + @pytest.mark.asyncio + async def test_the_profile_lands_in_the_run_dir(self, monkeypatch, tmp_path): + async def fake_videos(sec_user_id, count=20, *, cookie=""): + return [_video("111")] + + monkeypatch.setattr(douyin_api, "author_videos", fake_videos) + + await _collect(tmp_path, want_comments=False) + + assert self._profiles(tmp_path) == [_profile()] + + @pytest.mark.asyncio + async def test_the_profile_is_keyed_by_the_same_hash_as_the_works( + self, monkeypatch, tmp_path + ): + """**这条是关键。** 快照表的唯一键是 (任务, creator_hash, 轮次),而界面上是按 + 作品的 creator_hash 归组去查它的。两边只要差一个字符,粉丝数就永远查不出来 —— + 而且是静默的:表里有数据,界面上什么都没有。 + """ + + async def fake_videos(sec_user_id, count=20, *, cookie=""): + return [_video("111")] # 作品带的哈希是 "hash" + + # 资料接口自己算出来的是另一个值(比如它那边 uid 缺字段、只能拿 sec_uid 算)。 + async def off_hash_profile(sec_user_id, *, cookie=""): + return _profile(creator_hash="另一个哈希") + + monkeypatch.setattr(douyin_api, "author_videos", fake_videos) + monkeypatch.setattr(douyin_api, "author_profile", off_hash_profile) + + await _collect(tmp_path, want_comments=False) + + # 以作品为准:作品才是界面上的行,快照必须挂在能查到它的那个键上。 + assert self._profiles(tmp_path)[0]["creator_hash"] == "hash" + + @pytest.mark.asyncio + async def test_a_profile_without_an_identity_is_dropped(self, monkeypatch, tmp_path): + """哈希都算不出来的快照,落下去只会是一条谁也查不到的垃圾。""" + + async def fake_videos(sec_user_id, count=20, *, cookie=""): + return [] + + async def anonymous_profile(sec_user_id, *, cookie=""): + return _profile(creator_hash="") + + monkeypatch.setattr(douyin_api, "author_videos", fake_videos) + monkeypatch.setattr(douyin_api, "author_profile", anonymous_profile) + + result = await _collect(tmp_path, want_comments=False) + + assert self._profiles(tmp_path) == [] + assert any("身份标识" in error for error in result["errors"]) + + @pytest.mark.asyncio + async def test_a_failing_profile_does_not_lose_the_works(self, monkeypatch, tmp_path): + """附加信息拿不到,这一轮采到的作品不能跟着判成失败。""" + + async def fake_videos(sec_user_id, count=20, *, cookie=""): + return [_video("111")] + + async def broken_profile(sec_user_id, *, cookie=""): + raise douyin_api.DouyinApiError("资料接口抽风") + + monkeypatch.setattr(douyin_api, "author_videos", fake_videos) + monkeypatch.setattr(douyin_api, "author_profile", broken_profile) + + result = await _collect(tmp_path, want_comments=False) + + assert result["notes"] == 1 + assert self._profiles(tmp_path) == [] + assert any("资料失败" in error for error in result["errors"]) + + @pytest.mark.asyncio + async def test_a_work_target_never_asks_for_a_profile(self, monkeypatch, tmp_path): + """作品模式的目标是一件作品,没有「这个博主是谁」可问 —— 不该白发一个请求。""" + + asked = [] + + async def fake_detail(aweme_id, *, cookie=""): + return _video(aweme_id) + + async def recording_profile(sec_user_id, *, cookie=""): + asked.append(sec_user_id) + return _profile() + + monkeypatch.setattr(douyin_api, "video_detail", fake_detail) + monkeypatch.setattr(douyin_api, "author_profile", recording_profile) + + await _collect(tmp_path, mode="note", want_comments=False) + + assert asked == [] + # 文件仍然建出来(空的):事后翻 run 目录能看出「这次根本没问过」。 + assert self._profiles(tmp_path) == [] diff --git a/tests/test_monitor_creators.py b/tests/test_monitor_creators.py index f67a19d..05bbd9a 100644 --- a/tests/test_monitor_creators.py +++ b/tests/test_monitor_creators.py @@ -1,8 +1,9 @@ # -*- coding: utf-8 -*- -"""博主备注 —— 作品栏里那一位到底是谁。 +"""作品栏里的**博主**:备注(他到底是谁)与账号级指标(他现在多大)。 按 creator_hash 分组、显示 creator_name,两样都认不出人:一个是哈希,一个是平台昵称。 -备注是人自己起的名字。 +备注是人自己起的名字。账号级指标则是作品列表给不了的东西 —— 作品说的是"这条涨了多少赞", +粉丝数说的是"这个人整个账号在涨还是在掉"。 """ import httpx @@ -11,7 +12,12 @@ import pytest_asyncio from api.main import app from api.monitor import db as monitor_db -from api.monitor.models import MODE_CREATOR, MonitorNote, MonitorTask +from api.monitor.models import ( + MODE_CREATOR, + MonitorCreatorStat, + MonitorNote, + MonitorTask, +) CREATOR_HASH = "hash-a" NICKNAME = "张三" @@ -122,3 +128,93 @@ class TestCreatorAlias: await _set_alias(client, " 竞品A ") assert (await _notes(client))[0]["creator_alias"] == "竞品A" + + +async def _add_stat( + task_id: int, + run_id: int, + fans: int | None, + *, + total_favorited: int | None = 83000, + works: int | None = 42, + creator_hash: str = CREATOR_HASH, + captured_at: int = 1, +) -> None: + async with monitor_db.get_session() as session: + session.add( + MonitorCreatorStat( + task_id=task_id, + run_id=run_id, + creator_hash=creator_hash, + nickname=NICKNAME, + fans=fans, + total_favorited=total_favorited, + works_count=works, + following=7, + captured_at=captured_at, + ) + ) + + +class TestCreatorStats: + """账号级指标跟着作品一起返回 —— 界面上是按博主归组的,为了一个组头再发一轮请求 + 没道理。""" + + @pytest.mark.asyncio + async def test_without_snapshots_the_fields_are_null(self, client): + """null 而不是 0:0 会显示成「粉丝 0」,而事实是"还没采到"。""" + await _seed() + + note = (await _notes(client))[0] + + assert note["creator_fans"] is None + assert note["creator_total_favorited"] is None + assert note["creator_works"] is None + + @pytest.mark.asyncio + async def test_a_snapshot_shows_up_on_every_work_of_that_creator(self, client): + task_id = await _seed() + await _add_stat(task_id, run_id=1, fans=12000) + + for note in await _notes(client): + assert note["creator_fans"] == 12000 + assert note["creator_total_favorited"] == 83000 + assert note["creator_works"] == 42 + + @pytest.mark.asyncio + async def test_the_latest_run_wins(self, client): + """一轮一条,所以总会有好几条 —— 给界面的必须是最近那条。""" + task_id = await _seed() + await _add_stat(task_id, run_id=1, fans=12000) + await _add_stat(task_id, run_id=2, fans=12300) + + assert (await _notes(client))[0]["creator_fans"] == 12300 + + @pytest.mark.asyncio + async def test_another_platforms_snapshot_does_not_leak(self, client): + """两个平台上恰好同名同哈希的博主是两个人 —— 快照挂在任务上,不该串。""" + await _seed(platform="xhs") + dy_task_id = await _seed(platform="dy") + await _add_stat(dy_task_id, run_id=1, fans=999) + + assert (await _notes(client, "xhs"))[0]["creator_fans"] is None + assert (await _notes(client, "dy"))[0]["creator_fans"] == 999 + + @pytest.mark.asyncio + async def test_an_unparsed_count_stays_null(self, client): + """快照在,但某一项没解析出来 —— 那一项必须是 null,不能变成 0。 + + 和上一条的区别:那条是"根本没有快照",这条是"有快照、其中一项平台没给"。 + 界面上两种都该是「—」。 + """ + task_id = await _seed() + await _add_stat( + task_id, run_id=1, fans=None, total_favorited=None, works=None + ) + + note = (await _notes(client))[0] + + assert note["creator_fans"] is None + assert note["creator_works"] is None + # 快照本身是有的(有采集时间),只是值不知道 —— 前端要能分开这两件事。 + assert note["creator_stats_at"] is not None diff --git a/tests/test_monitor_ingest.py b/tests/test_monitor_ingest.py index 84bd864..91da5f7 100644 --- a/tests/test_monitor_ingest.py +++ b/tests/test_monitor_ingest.py @@ -54,6 +54,7 @@ from api.monitor.models import ( MODE_CREATOR, MonitorBase, MonitorComment, + MonitorCreatorStat, MonitorEvent, MonitorNote, MonitorNoteMetric, @@ -127,11 +128,15 @@ def _write_run_dir( notes: List[Dict[str, Any]], comments: Optional[List[Dict[str, Any]]] = None, subdir: str = "xhs", + profiles: Optional[List[Dict[str, Any]]] = None, ) -> Path: """Write a run's jsonl output in the crawler's own layout. ``subdir`` 是**爬虫**落盘的目录名,不是监控层的平台 id —— 抖音那边这两者不同 (平台 id 是 ``dy``、目录是 ``douyin``),所以必须能分开指定,否则测不出那个差异。 + + ``profiles`` 为 None 时**不写**这个文件(小红书那条路根本不产生它),给列表时写 + ——包括空列表,那是「问了但没问到」。 """ jsonl_dir = root / subdir / "jsonl" jsonl_dir.mkdir(parents=True, exist_ok=True) @@ -147,6 +152,12 @@ def _write_run_dir( "\n".join(json.dumps(c, ensure_ascii=False) for c in comments), encoding="utf-8", ) + if profiles is not None: + profile_file = jsonl_dir / "creator_profile_2026-01-01.jsonl" + profile_file.write_text( + "\n".join(json.dumps(p, ensure_ascii=False) for p in profiles), + encoding="utf-8", + ) return root @@ -829,3 +840,139 @@ class TestFailureDiagnosis: events = await _events(db, EVENT_RUN_FAILED) assert "account blocked" in events[0].title + + +# -------------------------------------------------------------------------- +# 博主账号级指标(粉丝 / 总获赞 / 作品数) +# -------------------------------------------------------------------------- + + +def _profile(creator_hash: str = "hash", **extra) -> Dict[str, Any]: + """``creator_profile_*.jsonl`` 里的一行 —— 形状由 douyin_api.author_profile 决定。""" + record: Dict[str, Any] = { + "creator_hash": creator_hash, + "nickname": "博主", + "unique_id": "abc", + "fans": 12000, + "total_favorited": 83000, + "works": 42, + "following": 7, + } + record.update(extra) + return record + + +class TestCreatorStatSnapshots: + """**账号级**指标和作品级指标是两回事:后者说"这条视频涨了多少赞",前者说 + "这个人整个账号的粉丝在涨还是在掉"。作品列表给不了后者,所以单独存一张表。 + """ + + async def _ingest(self, db, tmp_path, notes, profiles, 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, profiles=profiles) + result = await ingest_run(db, run, task, tmp_path) + return task, run, result + + @pytest.mark.asyncio + async def test_a_profile_becomes_a_snapshot(self, db, tmp_path): + task, run, _result = await self._ingest(db, tmp_path, [_dy_note("1")], [_profile()]) + + stat = await db.scalar(select(MonitorCreatorStat)) + + assert stat is not None + assert (stat.task_id, stat.run_id) == (task.id, run.id) + assert stat.creator_hash == "hash" + assert stat.nickname == "博主" + assert (stat.fans, stat.total_favorited, stat.works_count) == (12000, 83000, 42) + + @pytest.mark.asyncio + async def test_a_missing_count_stays_null_not_zero(self, db, tmp_path): + """0 是真实值(掉到零),null 是不知道。混起来趋势图就是在撒谎。""" + await self._ingest(db, tmp_path, [_dy_note("1")], [_profile(fans=None)]) + + stat = await db.scalar(select(MonitorCreatorStat)) + + assert stat.fans is None + # 同一个博主其它字段照常。 + assert stat.works_count == 42 + + @pytest.mark.asyncio + async def test_abbreviated_counts_are_parsed(self, db, tmp_path): + """走的是和作品指标同一个 parse_count —— 平台给你「1.2万」也得认。""" + await self._ingest(db, tmp_path, [_dy_note("1")], [_profile(fans="1.2万")]) + + assert (await db.scalar(select(MonitorCreatorStat))).fans == 12000 + + @pytest.mark.asyncio + async def test_the_same_creator_twice_in_one_run_yields_one_snapshot(self, db, tmp_path): + """一个任务可以配多个目标,退化路径下它们可能落在同一个博主身上。 + + 唯一键是 (task_id, creator_hash, run_id) —— 重复插入会撞键,把整轮炸掉。 + (和作品重复那次是同一类事故。) + """ + await self._ingest( + db, tmp_path, [_dy_note("1")], [_profile(), _profile(nickname="另一条")] + ) + + stats = list((await db.scalars(select(MonitorCreatorStat))).all()) + + assert len(stats) == 1 + + @pytest.mark.asyncio + async def test_a_profile_without_a_hash_is_skipped(self, db, tmp_path): + """哈希都算不出来,这条快照谁也查不到,落下去只是垃圾。""" + await self._ingest(db, tmp_path, [_dy_note("1")], [_profile(creator_hash="")]) + + assert (await db.scalars(select(MonitorCreatorStat))).all() == [] + + @pytest.mark.asyncio + async def test_no_profile_file_is_fine(self, db, tmp_path): + """小红书那条路(爬虫进程)根本不产生这个文件 —— 不能因此报错。""" + task = await _make_task(db) # xhs + run = await _make_run(db, task, started_at=1) + _write_run_dir(tmp_path, [_note("n1")], comments=[]) # 没有 profiles 参数 + + result = await ingest_run(db, run, task, tmp_path) + + assert result.status == RUN_SUCCESS + assert (await db.scalars(select(MonitorCreatorStat))).all() == [] + + @pytest.mark.asyncio + async def test_the_stats_are_kept_even_when_no_works_were_fetched(self, db, tmp_path): + """**这条是这里最值得留的一个。** + + 作品列表被风控挡住时,这一轮一条作品都拿不到、run 会被判成失败。但博主的粉丝数 + 并不因为这件事就不存在 —— 「粉丝还在涨,但新作品没在发现」恰恰是最该看见的时刻。 + 快照要是挂在「作品采到了」后面,就正好在最需要它的那一轮丢掉。 + """ + _task, run, result = await self._ingest(db, tmp_path, [], [_profile()]) + + assert result.status == RUN_PARTIAL # 一条作品都没有,这轮确实不算成功 + assert run.status == RUN_PARTIAL + stat = await db.scalar(select(MonitorCreatorStat)) + assert stat is not None and stat.fans == 12000 + + @pytest.mark.asyncio + async def test_each_run_adds_its_own_snapshot(self, db, tmp_path): + """趋势靠的就是这个:一条一轮,不要覆盖。""" + task = await _make_task(db, platform="dy") + first = await _make_run(db, task, started_at=1) + _write_run_dir(tmp_path, [_dy_note("1")], comments=[], subdir="douyin", + profiles=[_profile(fans=12000)]) + await ingest_run(db, first, task, tmp_path) + + second = await _make_run(db, task, started_at=2000) + _write_run_dir(tmp_path, [_dy_note("1")], comments=[], subdir="douyin", + profiles=[_profile(fans=12300)]) + await ingest_run(db, second, task, tmp_path) + + stats = list( + ( + await db.scalars( + select(MonitorCreatorStat).order_by(MonitorCreatorStat.run_id) + ) + ).all() + ) + + assert [s.fans for s in stats] == [12000, 12300] diff --git a/tests/test_monitor_note_alias.py b/tests/test_monitor_note_alias.py new file mode 100644 index 0000000..886c5a0 --- /dev/null +++ b/tests/test_monitor_note_alias.py @@ -0,0 +1,146 @@ +# -*- coding: utf-8 -*- +"""作品备注 —— 一个博主底下,哪几条是真正要盯的。 + +和博主备注(test_monitor_creators.py)是一对,但回答的不是同一个问题:博主备注回答 +「这个账号是谁」,作品备注回答「这条作品我要盯着」。一个博主底下常常只有一两件值得 +盯的作品,所以不能合成一条。 +""" + +import httpx +import pytest +import pytest_asyncio + +from api.main import app +from api.monitor import db as monitor_db +from api.monitor.models import MODE_CREATOR, MonitorNote, MonitorTask + +NOTE_ID = "note-a" +TITLE = "中秋哪儿都堵" + + +@pytest_asyncio.fixture +async def client(tmp_path): + monitor_db.set_sqlite_path(tmp_path / "monitor.db") + await monitor_db.init_db() + transport = httpx.ASGITransport(app=app) + async with httpx.AsyncClient(transport=transport, base_url="http://test") as http_client: + yield http_client + await monitor_db.dispose_engine() + + +async def _seed(platform: str = "xhs", task_name: str = "任务", note_id: str = NOTE_ID) -> int: + async with monitor_db.get_session() as session: + task = MonitorTask( + name=task_name, platform=platform, mode=MODE_CREATOR, enabled=True, + interval_minutes=60, max_notes_count=20, enable_comments=False, + max_comments_count=50, run_timeout_seconds=3600, + notify_enabled=False, created_at=0, updated_at=0, + ) + session.add(task) + await session.flush() + session.add( + MonitorNote( + task_id=task.id, note_id=note_id, title=TITLE, + note_url="", cover="", creator_hash="hash-a", + creator_name="张三", source_kind="", published_at=None, + first_seen_run_id=1, first_seen_at=0, + last_seen_run_id=1, last_seen_at=0, + ) + ) + return task.id + + +async def _notes(client, platform: str = "xhs"): + return (await client.get("/api/monitor/notes", params={"platform": platform})).json()["notes"] + + +async def _set_alias(client, alias: str, platform: str = "xhs", note_id: str = NOTE_ID): + return await client.put( + f"/api/monitor/notes/{note_id}", + params={"platform": platform}, + json={"alias": alias}, + ) + + +class TestNoteAlias: + @pytest.mark.asyncio + async def test_notes_start_without_a_remark(self, client): + await _seed() + + assert (await _notes(client))[0]["note_alias"] == "" + + @pytest.mark.asyncio + async def test_a_remark_comes_back_with_the_notes(self, client): + await _seed() + + response = await _set_alias(client, "重点") + + assert response.status_code == 200 + note = (await _notes(client))[0] + assert note["note_alias"] == "重点" + # 备注是**叠加**在标题之上的,不是替换 —— 标题仍然是这条作品本身。 + assert note["title"] == TITLE + + @pytest.mark.asyncio + async def test_a_remark_is_shared_across_tasks(self, client): + """同一件作品被两个任务都监控时,备注只该填一次。""" + await _seed(task_name="任务甲") + await _seed(task_name="任务乙") + await _set_alias(client, "重点") + + for note in await _notes(client): + assert note["note_alias"] == "重点" + + @pytest.mark.asyncio + async def test_a_remark_does_not_leak_to_another_platform(self, client): + """作品的 id 是平台各自的编号体系 —— 抖音的 123 和小红书的 123 是两条作品。""" + await _seed(platform="xhs") + await _seed(platform="dy") + + await _set_alias(client, "小红书那边的", platform="xhs") + + assert (await _notes(client, "xhs"))[0]["note_alias"] == "小红书那边的" + assert (await _notes(client, "dy"))[0]["note_alias"] == "" + + @pytest.mark.asyncio + async def test_a_remark_does_not_leak_to_another_work(self, client): + """钉住这里的**作用域**:键是 note_id。写错成按任务存的话,给一条起了备注, + 同一个博主底下的其它作品会跟着一起变 —— 那这个功能就没用了。""" + await _seed(note_id="note-a") + await _seed(task_name="另一个任务", note_id="note-b") + + await _set_alias(client, "重点", note_id="note-a") + + by_id = {note["note_id"]: note["note_alias"] for note in await _notes(client)} + assert by_id == {"note-a": "重点", "note-b": ""} + + @pytest.mark.asyncio + async def test_an_empty_remark_clears_it(self, client): + await _seed() + await _set_alias(client, "重点") + + await _set_alias(client, "") + + assert (await _notes(client))[0]["note_alias"] == "" + + @pytest.mark.asyncio + async def test_a_remark_is_trimmed(self, client): + await _seed() + + await _set_alias(client, " 重点 ") + + assert (await _notes(client))[0]["note_alias"] == "重点" + + @pytest.mark.asyncio + async def test_the_two_kinds_of_remark_stay_apart(self, client): + """博主备注和作品备注是两张表、两个键 —— 一个不该把另一个盖掉。""" + await _seed() + + await _set_alias(client, "重点") + await client.put( + "/api/monitor/creators/hash-a", params={"platform": "xhs"}, json={"alias": "竞品A"} + ) + + note = (await _notes(client))[0] + assert note["note_alias"] == "重点" + assert note["creator_alias"] == "竞品A" diff --git a/webui/src/components/monitor/NotesTable.tsx b/webui/src/components/monitor/NotesTable.tsx index a0f7a2f..3855c46 100644 --- a/webui/src/components/monitor/NotesTable.tsx +++ b/webui/src/components/monitor/NotesTable.tsx @@ -2,7 +2,7 @@ import { Fragment, useMemo, useRef, useState } from 'react' import { ChevronDown, ChevronRight, ExternalLink, Pencil, Users } from 'lucide-react' import { Badge } from '@/components/ui/badge' -import { useMonitorNotes, useSetCreatorAlias } from '@/hooks/useMonitor' +import { useMonitorNotes, useSetCreatorAlias, useSetNoteAlias } from '@/hooks/useMonitor' import { formatCount, formatDate, @@ -83,6 +83,32 @@ function sumMetrics(notes: MonitorNote[]) { // 多出来的那几列:展开箭头、作品、发布日期、首次发现、跳转链接;分组表头会跨掉整行。 const COLUMN_COUNT = METRIC_COLUMNS.length + 5 +/** + * 博主的**账号级**指标:粉丝 / 总获赞 / 作品数。 + * + * 作品列表给不了这个 —— 那几列说的是「这条作品涨了多少赞」,这里说的是「这个人整个 + * 账号在涨还是在掉」。**一个都没采到时整块不画**:画成「粉丝 0」比不画糟得多,那是 + * 一句假话。 + */ +function CreatorStats({ note }: { note: MonitorNote }) { + const parts = [ + note.creator_fans !== null && `粉丝 ${formatCount(note.creator_fans)}`, + note.creator_total_favorited !== null && `获赞 ${formatCount(note.creator_total_favorited)}`, + note.creator_works !== null && `${formatCount(note.creator_works)} 作品`, + ].filter(Boolean) as string[] + + if (parts.length === 0) return null + + return ( + + {parts.join(' · ')} + + ) +} + /** * 作品表,**按博主分组**。 * @@ -102,6 +128,11 @@ export function NotesTable({ taskId, onlyNew }: NotesTableProps) { // 否则它会把你刚放弃的内容存进去。 const aliasCancelled = useRef(false) const setAlias = useSetCreatorAlias() + // 正在改备注的那条**作品**。和博主备注是两个互不干扰的编辑态。 + const [editingNote, setEditingNote] = useState(null) + const [noteAliasDraft, setNoteAliasDraft] = useState('') + const noteAliasCancelled = useRef(false) + const setNoteAlias = useSetNoteAlias() const commitAlias = (creatorHash: string) => { if (aliasCancelled.current) { @@ -113,6 +144,16 @@ export function NotesTable({ taskId, onlyNew }: NotesTableProps) { setAlias.mutate({ creatorHash, alias: aliasDraft }) } + const commitNoteAlias = (noteId: string) => { + if (noteAliasCancelled.current) { + noteAliasCancelled.current = false + setEditingNote(null) + return + } + setEditingNote(null) + setNoteAlias.mutate({ noteId, alias: noteAliasDraft }) + } + const groups = useMemo(() => { const byCreator = new Map< string, @@ -227,6 +268,8 @@ export function NotesTable({ taskId, onlyNew }: NotesTableProps) { {group.notes.length} 篇 + {/* 账号级指标对组里每条作品都一样,取第一条即可。 */} + + + )} diff --git a/webui/src/hooks/useMonitor.ts b/webui/src/hooks/useMonitor.ts index 2abefd6..3013f1b 100644 --- a/webui/src/hooks/useMonitor.ts +++ b/webui/src/hooks/useMonitor.ts @@ -370,6 +370,23 @@ export function useSetCreatorAlias() { }) } +// --- 作品备注 -------------------------------------------------------------- + +export function useSetNoteAlias() { + const queryClient = useQueryClient() + const platform = usePlatformParam() + return useMutation({ + mutationFn: ({ noteId, alias }: { noteId: string; alias: string }) => + monitorApi.setNoteAlias(noteId, alias, platform), + onSuccess: () => { + toast.success('备注已保存') + queryClient.invalidateQueries({ queryKey: ['monitorNotes'] }) + queryClient.invalidateQueries({ queryKey: ['monitorCommentsGrouped'] }) + }, + onError: (error: Error) => toast.error(`备注保存失败:${error.message}`), + }) +} + // --- 上游更新检查 ----------------------------------------------------------------------------------------------------------------- export function useUpstreamStatus() { diff --git a/webui/src/lib/api.ts b/webui/src/lib/api.ts index ed7deaf..a7145d9 100644 --- a/webui/src/lib/api.ts +++ b/webui/src/lib/api.ts @@ -290,6 +290,13 @@ export const monitorApi = { { alias }, { params: { platform } }, ), + /** 给作品起备注(空串 = 清掉)。键是 note_id —— 和博主备注是两回事。 */ + setNoteAlias: (noteId: string, alias: string, platform?: string) => + api.put( + `/monitor/notes/${encodeURIComponent(noteId)}`, + { alias }, + { params: { platform } }, + ), /** * 立刻检查一次。服务端要等 fetch 跑完才回答,而 axios 默认 30 秒对此不够 —— * 一条卡住的 git fetch 能拖到两分钟,这里必须单独放长超时,否则会误报失败。 diff --git a/webui/src/types/monitor.ts b/webui/src/types/monitor.ts index 4c55e16..d3b3f62 100644 --- a/webui/src/types/monitor.ts +++ b/webui/src/types/monitor.ts @@ -102,6 +102,25 @@ export interface MonitorNote { * 更认不出。备注是唯一能把账号对上人的东西。空串表示没起过。 */ creator_alias: string + /** + * 人给**这条作品**起的备注。 + * + * 和 `creator_alias` 是两件事:那条回答「这个账号是谁」,这条回答「这条我要盯着」。 + * 一个博主底下常常只有一两件值得盯的作品,所以不能合并成一条。空串表示没起过。 + */ + note_alias: string + /** + * 博主**账号级**指标 —— 作品列表给不了的东西:作品说的是「这条涨了多少赞」, + * 它说的是「这个人整个账号在涨还是在掉」。 + * + * 三者都可能为 null:平台没采到(小红书那条路根本不产生它),或者某一项平台没给。 + * 为 null 时界面整块不画 —— 画成「粉丝 0」就是在撒谎。 + */ + creator_fans: number | null + creator_total_favorited: number | null + creator_works: number | null + /** 上述指标是哪一轮采到的。null = 从来没有过。 */ + creator_stats_at: number | null /** 作品的**发布**时间(爬虫侧:小红书 time、抖音 create_time)。可能与 * `first_seen_at` 差很远 —— 后者是「我们第一次看到它」的时间。平台没给时为 null。 */ published_at: number | null