From 9e13a7f686e653ad27d8911e2a9a889a4a6c3f9f Mon Sep 17 00:00:00 2001 From: butubb <1422726308@qq.com> Date: Sat, 10 Oct 2026 17:47:41 +0800 Subject: [PATCH] =?UTF-8?q?fix(monitor):=20=E6=8A=96=E9=9F=B3=E4=BB=BB?= =?UTF-8?q?=E5=8A=A1=E7=9A=84=20run=20=E6=B0=B8=E8=BF=9C=E5=81=9C=E5=9C=A8?= =?UTF-8?q?=E3=80=8C=E6=8E=92=E9=98=9F=E4=B8=AD=E3=80=8D=E2=80=94=E2=80=94?= =?UTF-8?q?=20=E6=88=91=E4=B8=8A=E4=B8=80=E7=89=88=E6=8A=8A=E7=8A=B6?= =?UTF-8?q?=E6=80=81=E6=A0=87=E8=AE=B0=E7=BC=A9=E8=BF=9B=E9=94=99=E4=BA=86?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 用户报的现象:任务一直显示「排队中」。查库确认有两批 run 卡在 pending(任务 7 的 44/45、 任务 15 的 62/63)。两个原因,一个是我上一版改坏的: 1) **`RUN_RUNNING` 被我缩进进了爬虫那条分支。** 抖音走的是另一条路,于是它**从不标记 「运行中」** —— 建完 pending 那一行就直接进采集,中途一旦出事(异常、进程被重启), 状态就永远停在 pending。这是我加平台分岔时把原本在两条路公共位置的一行挪进去了。 2) **`recover()` 只收 `running`,够不着 `pending`。** 那行是上一轮建的、后面的采集却 根本没机会开始(进程重启),它永远不会自己往前走。于是重启也救不回来,界面上就是 一个永远「排队中」的幽灵。现在 pending 一起收。 两处都补了测试:抖音路的 run 必须在**采集开始之前**就已经是 running(这条改回去就会 失败);recover 要把 pending 也标成 interrupted。 --- api/monitor/douyin_api.py | 49 ++++++++++++-------- api/monitor/runner.py | 15 +++--- api/monitor/scheduler.py | 16 ++++++- tests/test_monitor_runner.py | 82 +++++++++++++++++++++++++++++++++ tests/test_monitor_scheduler.py | 34 ++++++++++++++ 5 files changed, 168 insertions(+), 28 deletions(-) create mode 100644 tests/test_monitor_runner.py 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: