440 lines
15 KiB
Python
440 lines
15 KiB
Python
#!/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"
|
|
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())
|