#!/usr/bin/env python3 from __future__ import annotations import argparse import os import shlex import socket import subprocess import sys import time from dataclasses import dataclass from pathlib import Path from typing import Callable, Iterable, Optional import pymysql import yaml APP_DIR = Path(__file__).resolve().parents[1] if str(APP_DIR) not in sys.path: sys.path.insert(0, str(APP_DIR)) from main import get_action_timeout # noqa: E402 from src.core.config import load_config # noqa: E402 STATUS_RUNNING = 0 STATUS_SUCCESS = 1 STATUS_FAILED = 2 STATUS_SKIPPED = 3 OUTPUT_LIMIT = 5000 @dataclass class CrawlerTask: task_key: str task_name: str action: str cron_expression: str active: bool = True timeout: Optional[int] = None @classmethod def from_dict(cls, row: dict) -> "CrawlerTask": action = str(row.get("action") or row.get("params") or "").strip() return cls( task_key=str(row.get("key") or row.get("task_key") or action).strip(), task_name=str(row.get("name") or row.get("task_name") or action).strip(), action=action, cron_expression=str(row.get("cron") or row.get("cron_expression") or "* * * * *").strip(), active=bool(row.get("active", True)), timeout=int(row["timeout"]) if row.get("timeout") not in (None, "") else None, ) def resolved_timeout(self) -> int: return int(self.timeout or get_action_timeout(self.action)) @dataclass class CommandResult: exit_code: int output: str elapsed: float timed_out: bool = False @dataclass class TaskRunResult: log_id: int status: int output: str elapsed: float error_message: str = "" class TaskLock: def __init__(self, action: str, lock_dir: Path): safe_action = "".join(ch if ch.isalnum() or ch in ("-", "_") else "_" for ch in str(action or "all")) self.lock_dir = Path(lock_dir) self.path = self.lock_dir / f"{safe_action}.lock" self._fp = None self._locked = False def acquire(self) -> bool: self.lock_dir.mkdir(parents=True, exist_ok=True) self._fp = open(self.path, "a+b") try: self._fp.seek(0) if os.name == "nt": import msvcrt msvcrt.locking(self._fp.fileno(), msvcrt.LK_NBLCK, 1) else: import fcntl fcntl.flock(self._fp.fileno(), fcntl.LOCK_EX | fcntl.LOCK_NB) self._locked = True self._fp.seek(0) self._fp.truncate() self._fp.write(f"pid={os.getpid()} started_at={int(time.time())}\n".encode("utf-8")) self._fp.flush() return True except (BlockingIOError, OSError): self.release() return False def release(self) -> None: if not self._fp: return try: if self._locked: if os.name == "nt": import msvcrt self._fp.seek(0) msvcrt.locking(self._fp.fileno(), msvcrt.LK_UNLCK, 1) else: import fcntl fcntl.flock(self._fp.fileno(), fcntl.LOCK_UN) finally: self._locked = False self._fp.close() self._fp = None class TaskLogStore: def __init__(self): cfg = load_config().database self.prefix = cfg.prefix self.conn = pymysql.connect( host=cfg.host, port=cfg.port, user=cfg.username, password=cfg.password, database=cfg.database, charset=cfg.charset, autocommit=True, ) self.ensure_table() @property def table(self) -> str: return f"`{self.prefix}crawler_task_log`" def close(self) -> None: self.conn.close() def ensure_table(self) -> None: sql = f""" CREATE TABLE IF NOT EXISTS {self.table} ( `id` BIGINT UNSIGNED NOT NULL AUTO_INCREMENT, `task_key` VARCHAR(100) NOT NULL DEFAULT '' COMMENT '任务标识', `task_name` VARCHAR(100) NOT NULL DEFAULT '' COMMENT '任务名称', `action` VARCHAR(64) NOT NULL DEFAULT '' COMMENT '爬虫动作', `cron_expression` VARCHAR(64) NOT NULL DEFAULT '' COMMENT 'cron规则', `source` VARCHAR(32) NOT NULL DEFAULT 'docker' COMMENT '来源', `container_id` VARCHAR(128) NOT NULL DEFAULT '' COMMENT '容器ID', `status` TINYINT UNSIGNED NOT NULL DEFAULT 0 COMMENT '0执行中 1成功 2失败 3跳过', `trigger_type` TINYINT UNSIGNED NOT NULL DEFAULT 1 COMMENT '1自动 2手动', `output` TEXT NULL COMMENT '执行输出', `error_message` TEXT NULL COMMENT '错误信息', `elapsed` DECIMAL(10,2) NOT NULL DEFAULT 0.00 COMMENT '耗时秒', `started_at` INT UNSIGNED NOT NULL DEFAULT 0 COMMENT '开始时间', `finished_at` INT UNSIGNED NOT NULL DEFAULT 0 COMMENT '结束时间', `create_time` INT UNSIGNED NOT NULL DEFAULT 0, `update_time` INT UNSIGNED NOT NULL DEFAULT 0, PRIMARY KEY (`id`), KEY `idx_action_status` (`action`, `status`), KEY `idx_create_time` (`create_time`), KEY `idx_task_key` (`task_key`) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='Docker爬虫任务执行日志' """ with self.conn.cursor() as cur: cur.execute(sql) def create_running(self, task: CrawlerTask) -> int: now = int(time.time()) with self.conn.cursor() as cur: cur.execute( f""" INSERT INTO {self.table} (`task_key`, `task_name`, `action`, `cron_expression`, `source`, `container_id`, `status`, `trigger_type`, `output`, `error_message`, `elapsed`, `started_at`, `finished_at`, `create_time`, `update_time`) VALUES (%s, %s, %s, %s, 'docker', %s, %s, 1, '', '', 0, %s, 0, %s, %s) """, ( task.task_key, task.task_name, task.action, task.cron_expression, socket.gethostname(), STATUS_RUNNING, now, now, now, ), ) return int(cur.lastrowid) def update_progress(self, log_id: int, output: str, elapsed: float) -> None: with self.conn.cursor() as cur: cur.execute( f"UPDATE {self.table} SET output=%s, elapsed=%s, update_time=%s WHERE id=%s", (output[-OUTPUT_LIMIT:], round(elapsed, 2), int(time.time()), log_id), ) def finish(self, log_id: int, status: int, output: str, elapsed: float, error_message: str = "") -> None: now = int(time.time()) with self.conn.cursor() as cur: cur.execute( f""" UPDATE {self.table} SET status=%s, output=%s, error_message=%s, elapsed=%s, finished_at=%s, update_time=%s WHERE id=%s """, (status, output[-OUTPUT_LIMIT:], error_message[:2000], round(elapsed, 2), now, now, log_id), ) def load_tasks(path: Path) -> list[CrawlerTask]: raw = yaml.safe_load(Path(path).read_text(encoding="utf-8")) or {} rows = raw if isinstance(raw, list) else raw.get("tasks", []) return [CrawlerTask.from_dict(row) for row in rows] def render_crontab(tasks: Iterable[CrawlerTask], python_bin: str = "/usr/local/bin/python", app_dir: str = "/app") -> str: lines = [ "SHELL=/bin/bash", "BASH_ENV=/etc/sport-era-crawler/env.sh", "PATH=/usr/local/sbin:/usr/local/bin:/usr/sbin:/usr/bin:/sbin:/bin", "", ] for task in tasks: if not task.active: continue args = [ python_bin, "scripts/docker_task_runner.py", "run", task.action, "--task-key", task.task_key, "--task-name", task.task_name, "--cron-expression", task.cron_expression, ] if task.timeout: args.extend(["--timeout", str(task.timeout)]) cmd = " ".join(shlex.quote(str(arg)) for arg in args) lines.append(f"{task.cron_expression} root cd {shlex.quote(app_dir)} && {cmd}") lines.append("") return "\n".join(lines) def run_main_action(action: str, timeout: int, output_callback: Callable[[str], None]) -> CommandResult: start = time.monotonic() env = os.environ.copy() env["PYTHONUNBUFFERED"] = "1" env["DQD_RUNNING_UNDER_TASK_RUNNER"] = "1" proc = subprocess.Popen( [sys.executable, "main.py", action], cwd=str(APP_DIR), stdout=subprocess.PIPE, stderr=subprocess.STDOUT, text=True, encoding="utf-8", errors="replace", env=env, ) output_parts: list[str] = [] timed_out = False if os.name == "nt": try: stdout, _ = proc.communicate(timeout=timeout) stdout = stdout or "" output_parts.append(stdout) if stdout: output_callback(stdout) except subprocess.TimeoutExpired: timed_out = True proc.kill() stdout, _ = proc.communicate() stdout = stdout or "" output_parts.append(stdout) if stdout: output_callback(stdout) else: import selectors assert proc.stdout is not None selector = selectors.DefaultSelector() selector.register(proc.stdout, selectors.EVENT_READ) try: while proc.poll() is None: remaining = timeout - (time.monotonic() - start) if remaining <= 0: timed_out = True proc.kill() break for key, _ in selector.select(timeout=min(1.0, remaining)): line = key.fileobj.readline() if line: output_parts.append(line) output_callback(line) rest = proc.stdout.read() if rest: output_parts.append(rest) output_callback(rest) finally: selector.close() elapsed = round(time.monotonic() - start, 2) return CommandResult( exit_code=proc.returncode if proc.returncode is not None else -1, output="".join(output_parts), elapsed=elapsed, timed_out=timed_out, ) def run_task( task: CrawlerTask, store=None, lock_dir: Optional[Path] = None, command_runner: Callable[[str, int, Callable[[str], None]], CommandResult] = run_main_action, flush_interval: float = 5.0, ) -> TaskRunResult: if store is None: store = TaskLogStore() if hasattr(store, "ensure_table"): store.ensure_table() log_id = int(store.create_running(task)) start = time.monotonic() lock = TaskLock(task.action, Path(lock_dir or os.environ.get("DQD_LOCK_DIR") or APP_DIR / "data" / "locks")) output_parts: list[str] = [] last_flush = 0.0 def append_output(text: str) -> None: nonlocal last_flush output_parts.append(text) now = time.monotonic() elapsed = round(now - start, 2) if now - last_flush >= flush_interval: store.update_progress(log_id, "".join(output_parts), elapsed) last_flush = now if not lock.acquire(): output = f"{task.action} 已有任务执行中,跳过本次" store.finish(log_id, STATUS_SKIPPED, output, 0, "") return TaskRunResult(log_id, STATUS_SKIPPED, output, 0, "") try: result = command_runner(task.action, task.resolved_timeout(), append_output) output = result.output or "".join(output_parts) or "执行完成" if result.timed_out: status = STATUS_FAILED error = f"{task.action} 执行超时,超过 {task.resolved_timeout()} 秒" elif result.exit_code == 0: status = STATUS_SUCCESS error = "" else: status = STATUS_FAILED error = f"{task.action} 执行失败,退出码 {result.exit_code}" store.finish(log_id, status, output, result.elapsed, error) return TaskRunResult(log_id, status, output, result.elapsed, error) except Exception as exc: elapsed = round(time.monotonic() - start, 2) output = "".join(output_parts) or str(exc) store.finish(log_id, STATUS_FAILED, output, elapsed, str(exc)) return TaskRunResult(log_id, STATUS_FAILED, output, elapsed, str(exc)) finally: lock.release() def build_task_from_args(args) -> CrawlerTask: if args.tasks: for task in load_tasks(Path(args.tasks)): if task.action == args.action: if args.timeout: task.timeout = args.timeout return task return CrawlerTask( task_key=args.task_key or args.action, task_name=args.task_name or args.action, action=args.action, cron_expression=args.cron_expression or "", active=True, timeout=args.timeout, ) def parse_args(argv: Optional[list[str]] = None): parser = argparse.ArgumentParser(description="Sport Era Docker crawler task runner") sub = parser.add_subparsers(dest="command", required=True) run = sub.add_parser("run") run.add_argument("action") run.add_argument("--task-key", default="") run.add_argument("--task-name", default="") run.add_argument("--cron-expression", default="") run.add_argument("--timeout", type=int, default=None) run.add_argument("--tasks", default="") render = sub.add_parser("render-cron") render.add_argument("--tasks", required=True) render.add_argument("--output", required=True) render.add_argument("--python-bin", default="/usr/local/bin/python") render.add_argument("--app-dir", default="/app") return parser.parse_args(argv) def main(argv: Optional[list[str]] = None) -> int: args = parse_args(argv) if args.command == "render-cron": tasks = load_tasks(Path(args.tasks)) content = render_crontab(tasks, python_bin=args.python_bin, app_dir=args.app_dir) output = Path(args.output) output.parent.mkdir(parents=True, exist_ok=True) output.write_text(content, encoding="utf-8") output.chmod(0o644) return 0 if args.command == "run": store = TaskLogStore() try: result = run_task(build_task_from_args(args), store=store) return 0 if result.status in (STATUS_SUCCESS, STATUS_SKIPPED) else 1 finally: store.close() return 1 if __name__ == "__main__": raise SystemExit(main())