diff --git a/api/monitor/service.py b/api/monitor/service.py index 8b5fb7b..f0b7b21 100644 --- a/api/monitor/service.py +++ b/api/monitor/service.py @@ -460,15 +460,75 @@ async def _latest_creator_stats( return latest +def _window_order(note: MonitorNote) -> tuple: + """开窗排序:按**发布时间**倒序。 + + 采集端 ``aweme/post`` 拿的就是最新 N 条,界面的窗口必须和它口径一致 —— 否则会出现 + 「显示着的这条其实已经采不到了」。 + + **时间不知道的排最前(永不隐藏)**:因为解析不出日期就让一条作品从列表里消失,比 + 多显示一条糟得多。 + """ + return (note.published_at is None, note.published_at or 0) + + +async def _within_creator_window( + session: AsyncSession, notes: List[MonitorNote] +) -> List[MonitorNote]: + """每个博主只留**最新 N 条**,N 取任务上的 ``max_notes_count``。 + + **库里一条都不删。** 作品栏是「我正在盯的窗口」,库是账本:掉出窗口的作品采集端就 + 再也取不到了(它取的是最新 N 条),指标会冻在最后一次而**它看上去和正在跟踪的 + 作品一模一样** —— 那才是真正会骗人的地方。趋势图、报表、导出读的仍然是全量。 + + 只对**博主模式**开窗:笔记模式的 ``max_notes_count`` 管的是「每条目标拉多少」, + 而它的目标本身就是一件作品,不存在「一个博主的第 21 条」。 + """ + task_ids = {note.task_id for note in notes} + tasks = { + row.id: row + for row in ( + await session.scalars( + select(MonitorTask).where(MonitorTask.id.in_(list(task_ids))) + ) + ).all() + } + + by_creator: Dict[tuple, List[MonitorNote]] = {} + for note in notes: + by_creator.setdefault((note.task_id, note.creator_hash), []).append(note) + + kept: List[MonitorNote] = [] + for (note_task_id, _creator_hash), group in by_creator.items(): + task = tasks.get(note_task_id) + if task is None or task.mode != MODE_CREATOR: + kept.extend(group) + continue + group.sort(key=_window_order, reverse=True) + kept.extend(group[: task.max_notes_count]) + return kept + + async def list_notes( session: AsyncSession, task_id: Optional[int] = None, only_new: bool = False, limit: int = 200, platform: Optional[str] = None, + *, + windowed: bool = False, ) -> List[Dict[str, Any]]: - """Tracked notes with their latest metrics and change vs the previous run.""" - query = select(MonitorNote).order_by(MonitorNote.last_seen_at.desc()).limit(limit) + """Tracked notes with their latest metrics and change vs the previous run. + + ``windowed`` 打开后只返回每个博主最新 ``max_notes_count`` 条(见 + ``_within_creator_window``)。**默认关闭**:导出要的是全量账本,而默认开启的话, + 忘了传这个参数的那条路会静默少数据。 + """ + query = select(MonitorNote).order_by(MonitorNote.last_seen_at.desc()) + # 开窗时**不能**在这里 limit:得先把每个博主的窗开出来,顺序反了的话,全局上限会把 + # 某个博主最新的那几条挤掉,剩下的还都是别人的。 + if not windowed: + query = query.limit(limit) if task_id is not None: query = query.where(MonitorNote.task_id == task_id) if platform is not None: @@ -481,6 +541,13 @@ async def list_notes( if not notes: return [] + if windowed: + notes = await _within_creator_window(session, notes) + # 开完窗再截全局上限。窗口内的作品本来就不多,这里的 200 只是兜底。 + notes = notes[:limit] + if not notes: + return [] + note_ids = [note.note_id for note in notes] # Fetch every snapshot for these notes in one go and gather the two most diff --git a/api/routers/monitor.py b/api/routers/monitor.py index f935712..d585afe 100644 --- a/api/routers/monitor.py +++ b/api/routers/monitor.py @@ -133,7 +133,11 @@ async def list_notes( ): async with get_session() as session: return { - "notes": await service.list_notes(session, task_id, only_new, limit, platform), + # windowed=True:作品栏是「正在盯的窗口」,每个博主只画最新 max_notes_count 条。 + # 库里的全量还在,趋势图/报表/导出读的都是它。 + "notes": await service.list_notes( + session, task_id, only_new, limit, platform, windowed=True + ), # 博主**单独给一份**,而不是让前端从作品里推。作品推不出「一条作品都没有的 # 博主」—— 那正是最该显示的一类(还在涨粉,只是最近没发)。 "creators": await service.list_creators(session, task_id, platform), diff --git a/tests/test_monitor_notes_window.py b/tests/test_monitor_notes_window.py new file mode 100644 index 0000000..04d7f99 --- /dev/null +++ b/tests/test_monitor_notes_window.py @@ -0,0 +1,190 @@ +# -*- coding: utf-8 -*- +"""作品栏的**窗口** —— 每个博主只显示最新 N 条。 + +采集端 ``aweme/post`` 拿的就是最新 N 条,掉出这个窗口的作品**再也采不到**,指标会冻在 +最后一次。而它看上去和正在跟踪的作品一模一样 —— 那才是会骗人的地方。所以界面按同一个 +窗口显示,库里则一条不删(趋势图、报表、导出读的是全量)。 +""" + +import csv +import io +from datetime import datetime + +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, MODE_NOTE, MonitorNote, MonitorTask + +CREATOR = "hash-a" + +# 固定的基准时刻,避免用例依赖「今天是几号」。 +BASE = int(datetime(2026, 1, 1, 12, 0).timestamp() * 1000) + + +@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_task(mode: str = MODE_CREATOR, cap: int = 2) -> int: + async with monitor_db.get_session() as session: + task = MonitorTask( + name="窗口测试", platform="xhs", mode=mode, enabled=True, + interval_minutes=60, max_notes_count=cap, 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() + return task.id + + +async def _add_note( + task_id: int, + note_id: str, + published_at, + last_seen_at: int, + creator_hash: str = CREATOR, +) -> None: + async with monitor_db.get_session() as session: + session.add( + MonitorNote( + task_id=task_id, note_id=note_id, title=f"作品{note_id}", + note_url="", cover="", creator_hash=creator_hash, + creator_name="博主", source_kind="", published_at=published_at, + first_seen_run_id=1, first_seen_at=BASE, + last_seen_run_id=1, last_seen_at=last_seen_at, + ) + ) + + +async def _notes(client, **params): + return (await client.get("/api/monitor/notes", params=params)).json()["notes"] + + +def _ids(notes): + return sorted(note["note_id"] for note in notes) + + +class TestCreatorWindow: + @pytest.mark.asyncio + async def test_only_the_newest_n_survive(self, client): + """三条作品、上限两条 —— 最早的 `a` 掉出窗口。""" + task_id = await _seed_task(cap=2) + await _add_note(task_id, "a", BASE, BASE) + await _add_note(task_id, "b", BASE + 1000, BASE + 1000) + await _add_note(task_id, "c", BASE + 2000, BASE + 2000) + + assert _ids(await _notes(client)) == ["b", "c"] + + @pytest.mark.asyncio + async def test_the_window_is_by_publish_time_not_by_last_seen(self, client): + """**这条是这个功能的关键。** + + 「最后一次见到」和「最新发布」是两回事:降级路径(作品列表被风控挡住)会把**所有** + 已知作品都刷一遍,于是每条作品的 last_seen_at 都变成最新 —— 按它开窗等于没开。 + 窗口必须按发布时间,那才是采集端 `aweme/post` 返回的排序。 + """ + task_id = await _seed_task(cap=2) + # a 发布最早,但「最后一次见到」最新。 + await _add_note(task_id, "a", BASE, BASE + 9000) + await _add_note(task_id, "b", BASE + 1000, BASE + 1000) + await _add_note(task_id, "c", BASE + 2000, BASE + 2000) + + assert _ids(await _notes(client)) == ["b", "c"] + + @pytest.mark.asyncio + async def test_a_work_with_no_publish_time_is_never_hidden(self, client): + """**解析不出日期不该让一条作品消失** —— 多显示一条远好过悄悄少一条。""" + task_id = await _seed_task(cap=2) + await _add_note(task_id, "undated", None, BASE) + await _add_note(task_id, "b", BASE + 1000, BASE + 1000) + await _add_note(task_id, "c", BASE + 2000, BASE + 2000) + + assert "undated" in _ids(await _notes(client)) + + @pytest.mark.asyncio + async def test_each_creator_gets_its_own_window(self, client): + """上限是「每个博主 N 条」,不是「整个任务 N 条」。""" + task_id = await _seed_task(cap=1) + await _add_note(task_id, "a1", BASE, BASE, creator_hash="hash-a") + await _add_note(task_id, "a2", BASE + 1000, BASE + 1000, creator_hash="hash-a") + await _add_note(task_id, "b1", BASE, BASE, creator_hash="hash-b") + + assert _ids(await _notes(client)) == ["a2", "b1"] + + @pytest.mark.asyncio + async def test_a_note_mode_task_is_not_windowed(self, client): + """笔记模式的 max_notes_count 管的是「每条目标拉多少」,它的目标本身就是一件作品。""" + task_id = await _seed_task(mode=MODE_NOTE, cap=1) + await _add_note(task_id, "a", BASE, BASE) + await _add_note(task_id, "b", BASE + 1000, BASE + 1000) + await _add_note(task_id, "c", BASE + 2000, BASE + 2000) + + assert _ids(await _notes(client)) == ["a", "b", "c"] + + +class TestTheDatabaseKeepsEverything: + @pytest.mark.asyncio + async def test_the_window_hides_it_but_does_not_delete_it(self, client): + """界面是窗口,库是账本 —— 掉出窗口那条必须还在,否则趋势图和报表会缺历史。""" + from sqlalchemy import select + from sqlalchemy.ext.asyncio import AsyncSession + + from api.monitor.db import get_engine + + task_id = await _seed_task(cap=1) + await _add_note(task_id, "a", BASE, BASE) + await _add_note(task_id, "b", BASE + 1000, BASE + 1000) + + assert _ids(await _notes(client)) == ["b"] + + from sqlalchemy.ext.asyncio import async_sessionmaker + + factory = async_sessionmaker( + get_engine(), class_=AsyncSession, expire_on_commit=False + ) + async with factory() as session: + stored = list((await session.scalars(select(MonitorNote.note_id))).all()) + + assert sorted(stored) == ["a", "b"] + + @pytest.mark.asyncio + async def test_the_export_still_gets_everything(self, client): + """导出是**全量账本**,不是屏幕上那一屏 —— 界面上少显示几条是故意的, + 但导出的数据少几条就是在丢东西了。""" + task_id = await _seed_task(cap=1) + await _add_note(task_id, "a", BASE, BASE) + await _add_note(task_id, "b", BASE + 1000, BASE + 1000) + + response = await client.get( + "/api/monitor/export", params={"kind": "notes", "format": "csv"} + ) + rows = list(csv.DictReader(io.StringIO(response.content.decode("utf-8-sig")))) + + assert sorted(row["作品ID"] for row in rows) == ["a", "b"] + + +class TestTheGroupHeaderCounts: + @pytest.mark.asyncio + async def test_the_creator_list_reports_the_full_count(self, client): + """组头显示「显示了几篇 / 库里一共几篇」,两个数不一样是**对的**。 + + 屏幕上是开过窗的(这个博主最新 N 条),而 ``note_count`` 是库里这个博主的全部 + 作品数 —— 组头那行提示说的正是这件事。 + """ + task_id = await _seed_task(cap=1) + await _add_note(task_id, "a", BASE, BASE) + await _add_note(task_id, "b", BASE + 1000, BASE + 1000) + + creators = (await client.get("/api/monitor/notes", params={})).json()["creators"] + + assert creators[0]["note_count"] == 2