# -*- 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/services/crawler_manager.py # GitHub: https://github.com/NanmiCoder # Licensed under NON-COMMERCIAL LEARNING LICENSE 1.1 # # 声明:本代码仅供学习和研究目的使用。使用者应遵守以下原则: # 1. 不得用于任何商业用途。 # 2. 使用时应遵守目标平台的使用条款和robots.txt规则。 # 3. 不得进行大规模爬取或对平台造成运营干扰。 # 4. 应合理控制请求频率,避免给目标平台带来不必要的负担。 # 5. 不得用于任何非法或不当的用途。 # # 详细许可条款请参阅项目根目录下的LICENSE文件。 # 使用本代码即表示您同意遵守上述原则和LICENSE中的所有条款。 import asyncio import subprocess import signal import os from collections import deque from typing import Deque, Optional, List from datetime import datetime from pathlib import Path from ..schemas import CrawlerStartRequest, LogEntry from .interpreter import resolve_python_cmd # 留住多少行爬虫输出,供 run_and_wait 的调用方诊断失败原因。 # 子进程的输出本来只流向日志 WebSocket,监控层只看得到退出码 —— 于是「退出码 1」 # 成了运行历史里唯一的信息,真正的报错(比如抖音的 `DataFetchError: account blocked`) # 谁也看不到。留个尾巴,让失败原因能被写进 run.error_message。 OUTPUT_TAIL_LINES = 80 class CrawlerManager: """Crawler process manager""" def __init__(self): self._lock = asyncio.Lock() self.process: Optional[subprocess.Popen] = None self.status = "idle" self.started_at: Optional[datetime] = None self.current_config: Optional[CrawlerStartRequest] = None self._log_id = 0 self._logs: List[LogEntry] = [] self._read_task: Optional[asyncio.Task] = None # Project root directory self._project_root = Path(__file__).parent.parent.parent # Log queue - for pushing to WebSocket self._log_queue: Optional[asyncio.Queue] = None # Completion signalling for run_and_wait(). Polling `status` is unreliable # because stop() also resets it to "idle", and `self.process` gets replaced # by any concurrent start(), so waiters need an explicit event instead. self._done: asyncio.Event = asyncio.Event() self.last_exit_code: Optional[int] = None # 本次运行输出的末尾若干行。见 OUTPUT_TAIL_LINES。 self._output_tail: Deque[str] = deque(maxlen=OUTPUT_TAIL_LINES) def get_output_tail(self) -> List[str]: """最近一次运行的输出尾巴(最早的排前面)。 只在 run_and_wait() 返回之后读才有意义 —— 它等到读输出的任务收尾才唤醒。 """ return list(self._output_tail) @property def logs(self) -> List[LogEntry]: return self._logs def get_log_queue(self) -> asyncio.Queue: """Get or create log queue""" if self._log_queue is None: self._log_queue = asyncio.Queue() return self._log_queue def is_busy(self) -> bool: """Whether a crawler process is currently alive. This is the authoritative busy check -- `status` is a lagging indicator that manual stop() also resets. """ return self.process is not None and self.process.poll() is None async def run_and_wait( self, config: CrawlerStartRequest, extra_args: Optional[List[str]] = None, timeout: Optional[float] = None, ) -> int: """Start a crawler run and block until it exits, returning the exit code. Used by the monitor scheduler. Returns a negative value if the run was killed by `timeout` or if the process could not be started at all. """ started = await self.start(config, extra_args=extra_args) if not started: return -1 # Capture the process we just launched: a concurrent start() would # replace self.process, so poll this reference rather than the attribute. proc = self.process if proc is None: return -1 try: await asyncio.wait_for(self._done.wait(), timeout=timeout) except asyncio.TimeoutError: await self.stop() return -1 return self.last_exit_code if self.last_exit_code is not None else -1 def _create_log_entry(self, message: str, level: str = "info") -> LogEntry: """Create log entry""" self._log_id += 1 entry = LogEntry( id=self._log_id, timestamp=datetime.now().strftime("%H:%M:%S"), level=level, message=message ) self._logs.append(entry) # Keep last 500 logs if len(self._logs) > 500: self._logs = self._logs[-500:] return entry async def _push_log(self, entry: LogEntry): """Push log to queue""" # 这里是所有输出的唯一出口(读循环、收尾、以及管理器自己的提示都走它), # 所以尾巴挂在这儿最省事,也不会漏。 self._output_tail.append(entry.message) if self._log_queue is not None: try: self._log_queue.put_nowait(entry) except asyncio.QueueFull: pass def _parse_log_level(self, line: str) -> str: """Parse log level""" line_upper = line.upper() if "ERROR" in line_upper or "FAILED" in line_upper: return "error" elif "WARNING" in line_upper or "WARN" in line_upper: return "warning" elif "SUCCESS" in line_upper or "完成" in line or "成功" in line: return "success" elif "DEBUG" in line_upper: return "debug" return "info" async def start( self, config: CrawlerStartRequest, extra_args: Optional[List[str]] = None, ) -> bool: """Start crawler process""" async with self._lock: if self.process and self.process.poll() is None: return False # Clear old logs self._logs = [] self._log_id = 0 # Reset completion signalling for this run self._done.clear() self.last_exit_code = None # 尾巴只属于本次运行,否则上一轮的报错会混进这一轮的诊断里。 self._output_tail.clear() # Clear pending queue (don't replace object to avoid WebSocket broadcast coroutine holding old queue reference) if self._log_queue is None: self._log_queue = asyncio.Queue() else: try: while True: self._log_queue.get_nowait() except asyncio.QueueEmpty: pass # Build command line arguments cmd = self._build_command(config, extra_args=extra_args) # Log start information entry = self._create_log_entry(f"Starting crawler: {' '.join(cmd)}", "info") await self._push_log(entry) try: # Start subprocess self.process = subprocess.Popen( cmd, stdout=subprocess.PIPE, stderr=subprocess.STDOUT, text=True, encoding='utf-8', bufsize=1, cwd=str(self._project_root), env={**os.environ, "PYTHONUNBUFFERED": "1"} ) self.status = "running" self.started_at = datetime.now() self.current_config = config entry = self._create_log_entry( f"Crawler started on platform: {config.platform.value}, type: {config.crawler_type.value}", "success" ) await self._push_log(entry) # Start log reading task self._read_task = asyncio.create_task(self._read_output()) return True except Exception as e: self.status = "error" entry = self._create_log_entry(f"Failed to start crawler: {str(e)}", "error") await self._push_log(entry) return False async def stop(self) -> bool: """Stop crawler process""" async with self._lock: if not self.process or self.process.poll() is not None: return False self.status = "stopping" entry = self._create_log_entry("Sending SIGTERM to crawler process...", "warning") await self._push_log(entry) try: self.process.send_signal(signal.SIGTERM) # Wait for graceful exit (up to 15 seconds) for _ in range(30): if self.process.poll() is not None: break await asyncio.sleep(0.5) # If still not exited, force kill if self.process.poll() is None: entry = self._create_log_entry("Process not responding, sending SIGKILL...", "warning") await self._push_log(entry) self.process.kill() entry = self._create_log_entry("Crawler process terminated", "info") await self._push_log(entry) except Exception as e: entry = self._create_log_entry(f"Error stopping crawler: {str(e)}", "error") await self._push_log(entry) self.status = "idle" self.current_config = None # Cancel log reading task if self._read_task: self._read_task.cancel() self._read_task = None return True def get_status(self) -> dict: """Get current status""" return { "status": self.status, "platform": self.current_config.platform.value if self.current_config else None, "crawler_type": self.current_config.crawler_type.value if self.current_config else None, "started_at": self.started_at.isoformat() if self.started_at else None, "error_message": None } def _build_command( self, config: CrawlerStartRequest, extra_args: Optional[List[str]] = None, ) -> list: """Build main.py command line arguments""" cmd = [*resolve_python_cmd(), "main.py"] cmd.extend(["--platform", config.platform.value]) cmd.extend(["--lt", config.login_type.value]) cmd.extend(["--type", config.crawler_type.value]) cmd.extend(["--save_data_option", config.save_option.value]) # Pass different arguments based on crawler type if config.crawler_type.value == "search" and config.keywords: cmd.extend(["--keywords", config.keywords]) elif config.crawler_type.value == "detail" and config.specified_ids: cmd.extend(["--specified_id", config.specified_ids]) elif config.crawler_type.value == "creator" and config.creator_ids: cmd.extend(["--creator_id", config.creator_ids]) if config.start_page != 1: cmd.extend(["--start", str(config.start_page)]) cmd.extend(["--get_comment", "true" if config.enable_comments else "false"]) cmd.extend(["--get_sub_comment", "true" if config.enable_sub_comments else "false"]) cmd.extend(["--get_media", "true" if config.enable_media else "false"]) if config.max_notes_count is not None: cmd.extend(["--crawler_max_notes_count", str(config.max_notes_count)]) if config.max_comments_count is not None: cmd.extend(["--max_comments_count_singlenotes", str(config.max_comments_count)]) # Each of these is only appended when explicitly set, so manual runs from # the Crawl tab keep exactly their previous behaviour. if config.save_data_path: cmd.extend(["--save_data_path", config.save_data_path]) if config.enable_cdp_mode is not None: cmd.extend(["--enable_cdp_mode", "true" if config.enable_cdp_mode else "false"]) if config.inject_all_cookies is not None: cmd.extend(["--inject_all_cookies", "true" if config.inject_all_cookies else "false"]) if config.save_login_state is not None: cmd.extend(["--save_login_state", "true" if config.save_login_state else "false"]) if config.max_concurrency_num is not None: cmd.extend(["--max_concurrency_num", str(config.max_concurrency_num)]) if config.crawler_max_sleep_sec is not None: cmd.extend(["--crawler_max_sleep_sec", str(config.crawler_max_sleep_sec)]) if config.enable_ip_proxy is not None: cmd.extend(["--enable_ip_proxy", "true" if config.enable_ip_proxy else "false"]) if config.ip_proxy_pool_count is not None: cmd.extend(["--ip_proxy_pool_count", str(config.ip_proxy_pool_count)]) if config.ip_proxy_provider_name: cmd.extend(["--ip_proxy_provider_name", config.ip_proxy_provider_name]) if config.static_proxy_url: cmd.extend(["--static_proxy_url", config.static_proxy_url]) # Prefer a cookie file over passing the cookie on the command line, where # it would be visible in the process list. if config.cookies_file: cmd.extend(["--cookies_file", config.cookies_file]) elif config.cookies: cmd.extend(["--cookies", config.cookies]) cmd.extend(["--headless", "true" if config.headless else "false"]) if extra_args: cmd.extend(extra_args) return cmd async def _read_output(self): """Asynchronously read process output""" loop = asyncio.get_event_loop() # Capture the process this reader was started for. self.process can be # replaced by a subsequent start(), which would otherwise make us read # the exit code of the wrong run. proc = self.process try: while proc and proc.poll() is None: # Read a line in thread pool line = await loop.run_in_executor( None, proc.stdout.readline ) if line: line = line.strip() if line: level = self._parse_log_level(line) entry = self._create_log_entry(line, level) await self._push_log(entry) # Read remaining output if proc and proc.stdout: remaining = await loop.run_in_executor( None, proc.stdout.read ) if remaining: for line in remaining.strip().split('\n'): if line.strip(): level = self._parse_log_level(line) entry = self._create_log_entry(line.strip(), level) await self._push_log(entry) # Process ended if self.status == "running": exit_code = proc.returncode if proc else -1 if exit_code == 0: entry = self._create_log_entry("Crawler completed successfully", "success") else: entry = self._create_log_entry(f"Crawler exited with code: {exit_code}", "warning") await self._push_log(entry) self.status = "idle" except asyncio.CancelledError: pass except Exception as e: entry = self._create_log_entry(f"Error reading output: {str(e)}", "error") await self._push_log(entry) finally: # Record the exit code and wake any run_and_wait() waiter. Runs in a # finally so a cancelled read task still releases the waiter. self.last_exit_code = proc.returncode if proc else None self._done.set() # Global singleton crawler_manager = CrawlerManager()