# -*- coding: utf-8 -*- # Copyright (c) 2025 relakkes@gmail.com # # This file is part of MediaCrawler project. # Repository: https://github.com/NanmiCoder/MediaCrawler/blob/main/api/monitor/runner.py # GitHub: https://github.com/NanmiCoder # Licensed under NON-COMMERCIAL LEARNING LICENSE 1.1 # # 声明:本代码仅供学习和研究目的使用。使用者应遵守以下原则: # 1. 不得用于任何商业用途。 # 2. 使用时应遵守目标平台的使用条款和robots.txt规则。 # 3. 不得进行大规模爬取或对平台造成运营干扰。 # 4. 应合理控制请求频率,避免给目标平台带来不必要的负担。 # 5. 不得用于任何非法或不当的用途。 # # 详细许可条款请参阅项目根目录下的LICENSE文件。 # 使用本代码即表示您同意遵守上述原则和LICENSE中的所有条款。 """Execute a single monitoring run: build the command, wait, then ingest. Runs reuse ``CrawlerManager`` so that monitor crawls share the existing single-subprocess guarantee and their logs stream to the existing Terminal component over the existing log WebSocket. """ import asyncio import os from pathlib import Path from typing import Iterable, List, Optional from tools.time_util import get_current_timestamp from ..schemas import ( CrawlerStartRequest, CrawlerTypeEnum, LoginTypeEnum, PlatformEnum, SaveDataOptionEnum, ) from ..services import crawler_manager from . import app_settings, notify from .db import get_session from .ingest import IngestResult, ingest_run from .models import ( MODE_CREATOR, RUN_FAILED, RUN_PENDING, RUN_RUNNING, RUN_TIMEOUT, MonitorRun, MonitorTarget, MonitorTask, ) from .settings import get_cookie, mark_cookie_ok PROJECT_ROOT = Path(__file__).parent.parent.parent MONITOR_RUNS_DIR = PROJECT_ROOT / "data" / "monitor_runs" # Monitor platform ids align with PlatformEnum's values, but mapping explicitly # beats relying on that coincidence. _PLATFORM_ENUM = { "xhs": PlatformEnum.XHS, "dy": PlatformEnum.DOUYIN, "ks": PlatformEnum.KUAISHOU, "bili": PlatformEnum.BILIBILI, "wb": PlatformEnum.WEIBO, "tieba": PlatformEnum.TIEBA, "zhihu": PlatformEnum.ZHIHU, } _XHS_WEB_BASE = "https://www.xiaohongshu.com" _CREATOR_PATH = "/user/profile" _NOTE_PATH = "/explore" # Timeout used when the caller does not care; tasks carry their own. DEFAULT_RUN_TIMEOUT_SECONDS = 3600 def build_target_url(value: str, kind: str) -> str: """Turn a stored target into a URL the crawler's parser accepts. Always emits a full URL rather than a bare id: the XHS parser accepts a bare 24-hex id only, so the URL form is the safer universal input. The ``xsec_token`` is appended when present but is deliberately optional -- it expires, and the id alone is what keeps a long-running task alive. """ path = _CREATOR_PATH if kind == MODE_CREATOR else _NOTE_PATH return f"{_XHS_WEB_BASE}{path}/{value}" def build_target_urls(mode: str, targets: Iterable[MonitorTarget]) -> List[str]: urls = [] for target in targets: url = build_target_url(target.external_id, target.kind) if target.xsec_token: url = f"{url}?xsec_token={target.xsec_token}" if target.xsec_source: url = f"{url}&xsec_source={target.xsec_source}" urls.append(url) return urls async def _strategy_settings(session, platform: str) -> dict: """Crawl-strategy and proxy settings for one platform. Per-platform because the values genuinely differ: what is a safe request interval on one site is a rate limit on another. Read per run rather than cached, so a change takes effect on the next scheduled run. """ return { "enable_sub_comments": bool( await app_settings.get_value(session, "enable_sub_comments", platform, False) ), "crawl_sleep_sec": int( await app_settings.get_value(session, "crawl_sleep_sec", platform, 2) ), "enable_ip_proxy": bool( await app_settings.get_value(session, "enable_ip_proxy", platform, False) ), "proxy_provider": await app_settings.get_value( session, "proxy_provider", platform, "kuaidaili" ), "proxy_pool_count": int( await app_settings.get_value(session, "proxy_pool_count", platform, 2) ), "static_proxy_url": await app_settings.get_value( session, "static_proxy_url", platform, "" ), } def _write_cookie_file(path: Path, cookie: str) -> None: path.parent.mkdir(parents=True, exist_ok=True) path.write_text(cookie, encoding="utf-8") def _remove_cookie_file(path: Path) -> None: """Best-effort removal; the cookie is a credential, do not leave it around.""" try: os.remove(path) except OSError: pass async def execute_task(task_id: int, trigger: str = "manual") -> IngestResult: """Run one monitoring cycle for ``task_id`` and ingest its output. Split into three phases with separate short-lived DB sessions so no transaction is held open across the multi-minute subprocess run. """ # --- Phase 1: book the run and work out where its output goes ------------- async with get_session() as session: task = await session.get(MonitorTask, task_id) if task is None: raise ValueError(f"Monitor task {task_id} not found") targets = [target for target in task.targets if target.enabled] if not targets: raise ValueError(f"Monitor task {task_id} has no enabled targets") platform = task.platform urls = build_target_urls(task.mode, targets) cookie = await get_cookie(session, platform) strategy = await _strategy_settings(session, platform) # System-wide switch. On a headless server the crawler must attach to the # Chrome already listening on the debug port -- that browser is where the # operator scanned the login QR, so its profile is the login. Default False # keeps desktop runs launching a private browser exactly as before. cdp_enabled = await app_settings.get_value(session, "cdp_enabled", fallback=False) run = MonitorRun( task_id=task.id, trigger=trigger, status=RUN_PENDING, phase=task.mode, save_data_path="", queued_at=get_current_timestamp(), not_before=0, max_comments_count=task.max_comments_count if task.enable_comments else 0, ) session.add(run) await session.flush() run_id = run.id out_dir = MONITOR_RUNS_DIR / str(task.id) / str(run_id) run.save_data_path = str(out_dir) # Snapshot the values the subprocess needs; `task` is detached after commit. mode = task.mode enable_comments = task.enable_comments max_notes_count = task.max_notes_count max_comments_count = task.max_comments_count timeout_seconds = task.run_timeout_seconds # --- Phase 2: run the crawler outside any transaction --------------------- cookie_file = out_dir / ".cookies" _write_cookie_file(cookie_file, cookie) request = CrawlerStartRequest( platform=_PLATFORM_ENUM[platform], login_type=LoginTypeEnum.COOKIE, crawler_type=CrawlerTypeEnum.CREATOR if mode == MODE_CREATOR else CrawlerTypeEnum.DETAIL, creator_ids=",".join(urls) if mode == MODE_CREATOR else "", specified_ids=",".join(urls) if mode != MODE_CREATOR else "", start_page=1, enable_comments=enable_comments, enable_sub_comments=strategy["enable_sub_comments"], enable_media=False, save_option=SaveDataOptionEnum.JSONL, cookies="", headless=True, max_notes_count=max_notes_count, max_comments_count=max_comments_count, # Isolate this run's output: the crawler names files by date only, so # otherwise same-day runs would append into one shared file. save_data_path=str(out_dir), # Attach to the browser already running on CDP_DEBUG_PORT when the # operator enabled it; otherwise launch a private, throwaway browser. enable_cdp_mode=cdp_enabled, # Only injecting web_session is not enough to sign requests from a cold # browser profile. inject_all_cookies=True, save_login_state=True, cookies_file=str(cookie_file), max_concurrency_num=1, # Strategy + proxy, surfaced on the Settings page. crawler_max_sleep_sec=strategy["crawl_sleep_sec"], enable_ip_proxy=strategy["enable_ip_proxy"], ip_proxy_pool_count=strategy["proxy_pool_count"], ip_proxy_provider_name=strategy["proxy_provider"], 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: _remove_cookie_file(cookie_file) # --- Phase 3: ingest ------------------------------------------------------ async with get_session() as session: run = await session.get(MonitorRun, run_id) task = await session.get(MonitorTask, task_id) if run is None or task is None: raise ValueError(f"Run {run_id} or task {task_id} vanished during execution") if exit_code == -1 and not (out_dir / "xhs").exists(): # run_and_wait returns -1 when the process could not start or timed out. run.status = RUN_TIMEOUT run.finished_at = get_current_timestamp() run.exit_code = exit_code run.error_message = "Run was killed by timeout or failed to start" result = IngestResult(status=RUN_TIMEOUT, error=run.error_message) else: run.exit_code = exit_code run.finished_at = get_current_timestamp() result = await ingest_run(session, run, task, out_dir) # A run that authenticated fine is the only useful signal that the # stored cookie still works. if result.notes_fetched > 0: await mark_cookie_ok(session, task.platform) task.last_run_at = run.finished_at task.last_status = result.status task.last_error = result.error # --- Phase 4: notify ------------------------------------------------------ # Runs after the ingest transaction has committed, in its own session. A push # failure must never roll back collected data, and notify_run() swallows its # own errors for the same reason. async with get_session() as session: task = await session.get(MonitorTask, task_id) run = await session.get(MonitorRun, run_id) if task is not None and run is not None: await notify.notify_run(session, task, run) return result