diff --git a/api/monitor/douyin_api.py b/api/monitor/douyin_api.py index 76edafe..851ae40 100644 --- a/api/monitor/douyin_api.py +++ b/api/monitor/douyin_api.py @@ -375,6 +375,31 @@ async def author_profile(sec_user_id: str, *, cookie: str = "") -> Dict[str, Any } +def normalize_comment(comment: Dict[str, Any], aweme_id: str) -> Dict[str, Any]: + """把接口返回的一条评论,翻译成 store 落盘的那套键名。 + + HTTP 路线和页面路线共用它 —— 同一套键名,ingest 才不用关心数据是怎么来的。 + 刻意**不带** ``sub_comment_count`` / ``parent_comment_id`` 的猜测值:接口给了就用, + 没给就留空,不编。 + """ + user = comment.get("user") or {} + return { + "comment_id": str(comment.get("cid") or ""), + "aweme_id": aweme_id, + "content": comment.get("text") or "", + "nickname": user.get("nickname") or "", + "creator_hash": anonymize_user_id( + str(user.get("uid") or user.get("sec_uid") or "") + ), + # 同为秒;adapters 会换算。 + "create_time": _as_int(comment.get("create_time")), + "like_count": str(_as_int(comment.get("digg_count"))), + "sub_comment_count": str(_as_int(comment.get("reply_comment_total"))), + # 顶层评论在抖音里是 "0";adapters.parent_comment_id 会归一成空串。 + "parent_comment_id": str(comment.get("reply_id") or "0"), + } + + async def video_comments( aweme_id: str, count: int = 20, *, cookie: str = "" ) -> List[Dict[str, Any]]: @@ -399,26 +424,10 @@ async def video_comments( identity, ) - records = [] - for comment in payload.get("comments") or []: - user = comment.get("user") or {} - records.append( - { - "comment_id": str(comment.get("cid") or ""), - "aweme_id": aweme_id, - "content": comment.get("text") or "", - "nickname": user.get("nickname") or "", - "creator_hash": anonymize_user_id( - str(user.get("uid") or user.get("sec_uid") or "") - ), - # 同为秒;adapters 会换算。 - "create_time": _as_int(comment.get("create_time")), - "like_count": str(_as_int(comment.get("digg_count"))), - "sub_comment_count": str(_as_int(comment.get("reply_comment_total"))), - # 顶层评论在抖音里是 "0";adapters.parent_comment_id 会归一成空串。 - "parent_comment_id": str(comment.get("reply_id") or "0"), - } - ) + records = [ + normalize_comment(comment, aweme_id) + for comment in payload.get("comments") or [] + ] return records diff --git a/api/monitor/runner.py b/api/monitor/runner.py index 2d8ae02..ad0906e 100644 --- a/api/monitor/runner.py +++ b/api/monitor/runner.py @@ -210,6 +210,15 @@ async def execute_task(task_id: int, trigger: str = "manual") -> IngestResult: # 浏览器指纹参数(参数说 Mac + Chrome 125、UA 说 Linux + Chrome 155),网关回一个 # 200 + 空 body,然后被翻译成「account blocked」—— 看着像账号被封,其实什么都不是。 # 见 douyin_fetch / douyin_api。 + # 先落「运行中」—— **两条路都要**。原先这一行只写在爬虫那条分支里,于是抖音那条路上 + # run 一直停在 pending;一旦中途出事(异常、进程被重启),界面上就是一个永远 + # 「排队中」的幽灵,而且 recover() 也只清理 running、收不到它。 + async with get_session() as session: + run = await session.get(MonitorRun, run_id) + if run is not None: + run.status = RUN_RUNNING + run.started_at = get_current_timestamp() + in_process_tail: List[str] = [] if platform == adapters.PLATFORM_DY: fetched = await douyin_fetch.collect( @@ -274,12 +283,6 @@ async def execute_task(task_id: int, trigger: str = "manual") -> IngestResult: static_proxy_url=strategy["static_proxy_url"] or None, ) - async with get_session() as session: - run = await session.get(MonitorRun, run_id) - if run is not None: - run.status = RUN_RUNNING - run.started_at = get_current_timestamp() - try: exit_code = await crawler_manager.run_and_wait(request, timeout=timeout_seconds) finally: diff --git a/api/monitor/scheduler.py b/api/monitor/scheduler.py index 10e9dec..aaeddd2 100644 --- a/api/monitor/scheduler.py +++ b/api/monitor/scheduler.py @@ -51,7 +51,13 @@ from tools.time_util import get_current_timestamp from ..services import crawler_manager from . import app_settings, schedule, upstream from .db import get_session -from .models import MonitorRun, MonitorTask, RUN_INTERRUPTED, RUN_RUNNING +from .models import ( + RUN_INTERRUPTED, + RUN_PENDING, + RUN_RUNNING, + MonitorRun, + MonitorTask, +) from .runner import execute_task from .settings import get_cookie @@ -147,11 +153,17 @@ class MonitorScheduler: A run still marked ``running`` cannot be running -- its subprocess died with the previous process. Marking it interrupted stops it from blocking the UI as a phantom in-flight run. + + **``pending`` 同样是残留**:那一行是上一轮建的,可它后面的采集根本没机会开始 + (进程被重启,或者采集那条路抛了异常),所以它永远不会自己往前走。只清 running + 的话,它会永远挂在界面上显示「排队中」—— 用户看到的就是任务卡住了。 """ async with get_session() as session: stale = ( await session.scalars( - select(MonitorRun).where(MonitorRun.status == RUN_RUNNING) + select(MonitorRun).where( + MonitorRun.status.in_((RUN_RUNNING, RUN_PENDING)) + ) ) ).all() for run in stale: diff --git a/tests/test_monitor_runner.py b/tests/test_monitor_runner.py new file mode 100644 index 0000000..9ad933a --- /dev/null +++ b/tests/test_monitor_runner.py @@ -0,0 +1,82 @@ +# -*- coding: utf-8 -*- +"""monitor runner —— 尤其是抖音那条(不走子进程的)路的运行状态流转。 + +这条路的地位特殊:它不经过 ``crawler_manager``,所以爬虫那套「退出码 / 日志尾巴」的 +约定它一个都不沾。凡是写在那里面的东西,这条路都得单独有一份。 +""" + +import pytest +import pytest_asyncio +from sqlalchemy import select + +from tools.time_util import get_current_timestamp + +from api.monitor import db as monitor_db +from api.monitor import runner as runner_module +from api.monitor.models import ( + MODE_CREATOR, + RUN_RUNNING, + MonitorRun, + MonitorTarget, + MonitorTask, +) + + +@pytest_asyncio.fixture +async def db(tmp_path): + monitor_db.set_sqlite_path(tmp_path / "monitor.db") + await monitor_db.init_db() + yield monitor_db + await monitor_db.dispose_engine() + + +async def _make_douyin_task() -> int: + async with monitor_db.get_session() as session: + now = get_current_timestamp() + task = MonitorTask( + name="dy", platform="dy", mode=MODE_CREATOR, enabled=True, + interval_minutes=360, max_notes_count=20, enable_comments=False, + max_comments_count=20, run_timeout_seconds=3600, + notify_enabled=False, notify_failures=False, + created_at=now, updated_at=now, + ) + session.add(task) + await session.flush() + session.add( + MonitorTarget( + task_id=task.id, kind=MODE_CREATOR, external_id="MS4w-sec", + xsec_token="", xsec_source="", raw_value="MS4w-sec", + label="x", enabled=True, created_at=now, + ) + ) + return task.id + + +class TestDouyinRunStatus: + @pytest.mark.asyncio + async def test_the_run_is_marked_running_before_collecting(self, db, monkeypatch): + """**采集开始之前**,run 就必须已经是 running。 + + 这一行原先只写在爬虫那条分支里,于是抖音路上 run 一直停在 pending —— 一旦中途 + 出事(异常、或进程被重启),界面上就是一个永远「排队中」的幽灵,而且 recover() + 当时也只收 running、够不着它。 + """ + task_id = await _make_douyin_task() + seen = {} + + async def fake_collect(out_dir, **kwargs): + async with monitor_db.get_session() as session: + run = await session.scalar(select(MonitorRun).order_by(MonitorRun.id)) + seen["status"] = run.status + return { + "notes": 0, + "comments": 0, + "errors": ["故意失败"], + "jsonl_dir": str(out_dir), + } + + monkeypatch.setattr(runner_module.douyin_fetch, "collect", fake_collect) + + await runner_module.execute_task(task_id, trigger="manual") + + assert seen["status"] == RUN_RUNNING diff --git a/tests/test_monitor_scheduler.py b/tests/test_monitor_scheduler.py index a571362..466e997 100644 --- a/tests/test_monitor_scheduler.py +++ b/tests/test_monitor_scheduler.py @@ -30,6 +30,7 @@ from api.monitor.models import ( MonitorTarget, MonitorTask, RUN_INTERRUPTED, + RUN_PENDING, RUN_RUNNING, RUN_SUCCESS, ) @@ -279,6 +280,39 @@ class TestRecovery: assert run.status == RUN_INTERRUPTED assert run.finished_at is not None + @pytest.mark.asyncio + async def test_pending_runs_are_also_cleaned_up(self, db): + """挂在 ``pending`` 的 run 同样是残留,必须一起收。 + + 那一行是上一轮建的,可它后面的采集根本没机会开始(进程被重启,或采集那条路抛了 + 异常)。只清 ``running`` 的话,它会永远挂在界面上显示「排队中」—— + 用户看到的就是任务卡死了。 + """ + async with monitor_db.get_session() as session: + now = get_current_timestamp() + task = MonitorTask( + name="t", platform="dy", mode=MODE_CREATOR, enabled=True, + interval_minutes=60, max_notes_count=20, enable_comments=False, + max_comments_count=50, run_timeout_seconds=3600, + next_run_at=now, last_status="pending", created_at=now, updated_at=now, + ) + session.add(task) + await session.flush() + session.add( + MonitorRun( + task_id=task.id, trigger="manual", status=RUN_PENDING, + phase=MODE_CREATOR, save_data_path="", queued_at=now, not_before=0, + max_comments_count=50, + ) + ) + + await MonitorScheduler().recover() + + async with monitor_db.get_session() as session: + run = await session.scalar(select(MonitorRun)) + assert run.status == RUN_INTERRUPTED + assert run.finished_at is not None + @pytest.mark.asyncio async def test_completed_runs_are_left_alone(self, db): async with monitor_db.get_session() as session: