From dede5003a49cc9222eebf1bf9f83b98e96dbf0a2 Mon Sep 17 00:00:00 2001 From: hz4th_coder Date: Tue, 11 Aug 2026 12:44:18 +0800 Subject: [PATCH] =?UTF-8?q?=E9=80=9A=E7=94=A8=E7=88=AC=E8=99=AB=E7=B3=BB?= =?UTF-8?q?=E7=BB=9F=20v1.0.0:=20=E6=89=B9=E9=87=8F/=E5=AE=9A=E6=97=B6/?= =?UTF-8?q?=E8=87=AA=E5=8A=A8=E7=88=AC=E5=8F=96=20+=20Web=E7=AE=A1?= =?UTF-8?q?=E7=90=86=E7=95=8C=E9=9D=A2=20+=20=E5=9B=BE=E7=89=87/=E9=80=9A?= =?UTF-8?q?=E7=9F=A5/=E7=83=AD=E6=9B=B4=E6=96=B0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .gitignore | 6 + README.md | 92 +++++++++ app.py | 338 ++++++++++++++++++++++++++++++ engine.py | 432 +++++++++++++++++++++++++++++++++++++++ legacy/crawl.py | 168 +++++++++++++++ legacy/test_one.py | 48 +++++ legacy/urls.txt | 4 + notify.py | 43 ++++ requirements.txt | 3 + scheduler.py | 120 +++++++++++ start.sh | 54 +++++ static/app.js | 497 +++++++++++++++++++++++++++++++++++++++++++++ static/index.html | 143 +++++++++++++ static/style.css | 181 +++++++++++++++++ store.py | 122 +++++++++++ 15 files changed, 2251 insertions(+) create mode 100644 .gitignore create mode 100644 README.md create mode 100644 app.py create mode 100644 engine.py create mode 100644 legacy/crawl.py create mode 100644 legacy/test_one.py create mode 100644 legacy/urls.txt create mode 100644 notify.py create mode 100644 requirements.txt create mode 100644 scheduler.py create mode 100755 start.sh create mode 100644 static/app.js create mode 100644 static/index.html create mode 100644 static/style.css create mode 100644 store.py diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..174e51e --- /dev/null +++ b/.gitignore @@ -0,0 +1,6 @@ +__pycache__/ +*.pyc +data/*.json +data/cookies_*.json +logs/ +out/ diff --git a/README.md b/README.md new file mode 100644 index 0000000..b6a858a --- /dev/null +++ b/README.md @@ -0,0 +1,92 @@ +# 通用爬虫系统 (universal-crawler) + +基于原 tpu_crawler(Playwright + stealth + 系统 Chrome)改造的全能网页爬取管理系统。 +**可爬取任意网址**,自带 Web 管理界面,支持批量、定时、自动三种模式。 + +## 快速开始 + +```bash +./start.sh # 启动 (默认端口 16062) +./start.sh stop # 停止 +./start.sh restart # 重启 +./start.sh status # 查看状态 +``` + +启动后打开管理界面:`http://<服务器IP>:16062/` + +环境要求(本机已具备): +- Python 3.12 (`/home/hz1/miniconda3/envs/openclaw/bin/python`) +- flask / playwright / playwright-stealth +- 系统 Chrome `/usr/bin/google-chrome` + +## 功能总览 + +### 1. 前端管理界面 +- 任务卡片总览:状态、进度、统计、下次调度时间一目了然 +- 一键操作:开始 / 暂停 / 恢复 / 终止 / 编辑 / 删除 +- 任务详情:历次运行记录、结果明细表(HTML/TXT 在线预览)、图片缩略图、实时日志 +- 运行中的任务参数支持**热更新**(修改后从下一页起生效) + +### 2. 批量爬取模式 +一次粘贴多个网址(每行一个,`#` 注释),可配置: +- 项目名称、输出目录(默认 `out/<任务ID>`,可填绝对路径) +- 爬取间隔(随机秒数区间,防封 IP)、单页超时 +- 失败重试次数 / 重试间隔 +- 是否爬取图片(自动下载页面图片到 `xxx_img/` 目录) +- 完成后邮件通知(复用 send_email.py,默认发到 wlq@tphai.com) + +### 3. 定时爬取模式 +- **间隔调度**:每 N 分钟 / 小时 / 天 +- **cron 表达式**:5 段式(分 时 日 月 周,周 0/7=周日),如 `0 3 * * *` 每天凌晨 3 点 +- 可启用/停用调度,自动计算下次执行时间;到点自动开跑,跑完自动计算下一次 + +### 4. 自动爬取模式 +给定一个起始网址,系统自动从页面里发现链接、按规则筛选后 BFS 爬取: +- **包含规则**:只爬包含指定子串(或正则)的链接 +- **排除规则**:跳过匹配的链接(如 login、/tag/) +- **仅同域名**:限制在起始网站内 +- **最大页数 / 最大深度**:控制爬取规模 +- 其余参数(间隔、重试、图片、通知)同批量模式 + +## 输出文件 + +每个任务输出到独立目录(默认 `out/<任务ID>/`): +- `NNNN_<域名>_<时间戳>.html` — 完整网页 +- `NNNN_<域名>_<时间戳>.txt` — 提取的纯文本正文 +- `NNNN_<域名>_<时间戳>_img/` — 下载的图片(开启图片爬取时) + +## API 速查 + +| 方法 | 路径 | 说明 | +|---|---|---| +| GET | `/api/tasks` | 任务列表(含最新运行) | +| POST | `/api/tasks` | 创建任务 | +| GET/PUT/DELETE | `/api/tasks/` | 详情 / 修改 / 删除 | +| POST | `/api/tasks//start` | 开始爬取 | +| POST | `/api/tasks//pause` | 暂停 | +| POST | `/api/tasks//resume` | 恢复 | +| POST | `/api/tasks//stop` | 终止 | +| GET | `/api/runs/` | 运行详情(结果+日志) | +| GET | `/api/runs//logs?offset=N` | 增量日志 | +| GET | `/api/file?task_id=&path=` | 读取输出文件(HTML/TXT/图片) | + +## 目录结构 + +``` +universal-crawler/ +├── app.py # Flask 主服务 (端口 16062) +├── engine.py # 爬虫引擎 (批量/自动/图片/重试/暂停恢复) +├── scheduler.py # 定时调度器 (间隔 + cron) +├── store.py # 任务/运行记录持久化 (JSON) +├── notify.py # 邮件通知 +├── static/ # 前端管理界面 +├── data/ # 运行时数据 (任务、运行记录、cookie) +├── out/ # 默认爬取输出 +├── legacy/ # 原 tpu_crawler 脚本备份 +└── start.sh # 启动脚本 (PID 文件管理) +``` + +## 注意事项 +- 爬取前尊重目标站点 robots.txt / 使用条款,仅用于合法用途 +- 延迟别设太低(建议 ≥1s),高频访问会被封 IP +- 每个任务独立的 cookie 文件,互不干扰;被反爬拦截时删除 `data/cookies_<任务ID>.json` 重新验证 diff --git a/app.py b/app.py new file mode 100644 index 0000000..62c82dc --- /dev/null +++ b/app.py @@ -0,0 +1,338 @@ +# -*- coding: utf-8 -*- +""" +通用爬虫系统 - Web 管理后端 +启动: /home/hz1/miniconda3/envs/openclaw/bin/python app.py (默认端口 16062) +""" +import os +import threading +from datetime import datetime + +from flask import Flask, jsonify, request, send_file, send_from_directory + +import store +from engine import CrawlJob +from scheduler import Scheduler, cron_next, interval_delta + +HERE = os.path.dirname(os.path.abspath(__file__)) +PORT = int(os.environ.get("CRAWLER_PORT", "16062")) + +app = Flask(__name__, static_folder="static", static_url_path="") +JOBS = {} # task_id -> CrawlJob +JOBS_LOCK = threading.RLock() +STARTED = datetime.now() + +DEFAULT_CONFIG = { + "out_dir": "", # 留空 -> out/<任务ID> + "delay_min": 2, + "delay_max": 5, + "timeout": 60, + "crawl_images": False, + "retry_count": 2, + "retry_interval": 3, + "notify": False, + "notify_email": "wlq@tphai.com", +} + + +def now_str(): + return store.now_str() + + +def resolve_out_dir(task): + cfg = task.get("config", {}) or {} + if cfg.get("out_dir", "").strip(): + return cfg["out_dir"].strip() + return os.path.join(HERE, "out", task["id"]) + + +def make_run(task): + return { + "id": store.new_id("r"), + "task_id": task["id"], + "mode": task.get("mode", "batch"), + "status": "running", + "progress": {"done": 0, "total": 0, "current_url": "", "percent": 0}, + "started_at": now_str(), + "finished_at": "", + "stats": {"ok": 0, "fail": 0, "images": 0}, + "results": [], + "logs": [], + "out_dir": resolve_out_dir(task), + } + + +def persist_cb(task_id, run): + total = run["progress"].get("total") or 0 + done = run["progress"].get("done") or 0 + run["progress"]["percent"] = round(done * 100 / total) if total else 0 + store.save_run(task_id, run) + + +def start_run(task): + """为任务启动一次爬取, 返回 (run, error)""" + with JOBS_LOCK: + job = JOBS.get(task["id"]) + if job and job.is_running(): + return None, "该任务已有正在运行的爬取" + run = make_run(task) + store.add_run(task["id"], run) + job = CrawlJob(task, run, persist_cb) + JOBS[task["id"]] = job + job.start() + return run, None + + +def _stop_job(tid): + with JOBS_LOCK: + job = JOBS.get(tid) + if job and job.is_running(): + job.stop() + return True + return False + + +# ---------------- 页面 ---------------- + +@app.route("/") +def index(): + return send_from_directory(app.static_folder, "index.html") + + +# ---------------- API: 状态 ---------------- + +@app.route("/api/status") +def api_status(): + with JOBS_LOCK: + running = sum(1 for j in JOBS.values() if j.is_running()) + return jsonify({ + "app": "universal-crawler", + "version": "1.0.0", + "running": running, + "uptime": str(datetime.now() - STARTED).split(".")[0], + }) + + +# ---------------- API: 任务 ---------------- + +@app.route("/api/tasks", methods=["GET"]) +def api_tasks(): + tasks = store.load_tasks() + for t in tasks: + runs = store.get_runs(t["id"]) + t["latest_run"] = runs[-1] if runs else None + with JOBS_LOCK: + job = JOBS.get(t["id"]) + t["running"] = bool(job and job.is_running()) + tasks.sort(key=lambda x: x.get("created_at", ""), reverse=True) + return jsonify(tasks) + + +@app.route("/api/tasks", methods=["POST"]) +def api_create_task(): + body = request.get_json(force=True) or {} + name = (body.get("name") or "").strip() + mode = body.get("mode", "batch") + if not name: + return jsonify({"error": "项目名不能为空"}), 400 + if mode not in ("batch", "auto", "scheduled"): + return jsonify({"error": "无效的模式"}), 400 + + task = { + "id": store.new_id("t"), + "name": name, + "mode": mode, + "created_at": now_str(), + "updated_at": now_str(), + "config": {**DEFAULT_CONFIG, **(body.get("config") or {})}, + "urls": [u.strip() for u in (body.get("urls") or []) if u.strip()], + } + + if mode == "auto": + auto = dict(body.get("auto") or {}) + seed = (auto.get("seed_url") or "").strip() + if not seed: + return jsonify({"error": "自动模式需要填写起始网址"}), 400 + if not seed.startswith("http"): + seed = "https://" + seed + auto["seed_url"] = seed + task["auto"] = auto + elif mode == "scheduled": + sch = dict(body.get("schedule") or {}) + sch.setdefault("enabled", True) + sch.setdefault("type", "interval") + try: + if sch.get("type") == "cron": + expr = sch.get("cron") or "0 * * * *" + nn = cron_next(expr) + if not nn: + raise ValueError("cron 表达式在未来一年内无匹配时间") + sch["next_run"] = nn.strftime("%Y-%m-%d %H:%M:%S") + else: + sch.setdefault("interval_unit", "hours") + sch.setdefault("interval_value", 24) + sch["next_run"] = (datetime.now() + interval_delta(sch)).strftime("%Y-%m-%d %H:%M:%S") + except ValueError as e: + return jsonify({"error": f"调度配置错误: {e}"}), 400 + sch.setdefault("last_run", "") + sch.setdefault("runs_count", 0) + task["schedule"] = sch + if not task["urls"]: + return jsonify({"error": "定时任务需要网址列表"}), 400 + else: + if not task["urls"]: + return jsonify({"error": "请至少填写一个网址"}), 400 + + store.upsert_task(task) + return jsonify(task), 201 + + +@app.route("/api/tasks/", methods=["GET"]) +def api_task_detail(tid): + task = store.get_task(tid) + if not task: + return jsonify({"error": "任务不存在"}), 404 + task["runs"] = list(reversed(store.get_runs(tid))) + with JOBS_LOCK: + job = JOBS.get(tid) + task["running"] = bool(job and job.is_running()) + return jsonify(task) + + +@app.route("/api/tasks/", methods=["PUT"]) +def api_update_task(tid): + task = store.get_task(tid) + if not task: + return jsonify({"error": "任务不存在"}), 404 + body = request.get_json(force=True) or {} + with JOBS_LOCK: + job = JOBS.get(tid) + running = bool(job and job.is_running()) + if "name" in body and str(body["name"]).strip(): + task["name"] = str(body["name"]).strip() + if "urls" in body: + task["urls"] = [u.strip() for u in body["urls"] if u.strip()] + if "config" in body: + merged = {**task.get("config", {}), **body["config"]} + task["config"] = merged + if running: + job.update_config(body["config"]) # 运行中热更新 + if "auto" in body and task.get("mode") == "auto": + task["auto"] = {**task.get("auto", {}), **body["auto"]} + if "schedule" in body and task.get("mode") == "scheduled": + sch = {**task.get("schedule", {}), **body["schedule"]} + try: + if sch.get("type") == "cron": + nn = cron_next(sch.get("cron") or "0 * * * *") + if not nn: + raise ValueError("cron 表达式在未来一年内无匹配时间") + sch["next_run"] = nn.strftime("%Y-%m-%d %H:%M:%S") + else: + sch["next_run"] = (datetime.now() + interval_delta(sch)).strftime("%Y-%m-%d %H:%M:%S") + except ValueError as e: + return jsonify({"error": f"调度配置错误: {e}"}), 400 + task["schedule"] = sch + task["updated_at"] = now_str() + store.upsert_task(task) + return jsonify(task) + + +@app.route("/api/tasks/", methods=["DELETE"]) +def api_delete_task(tid): + _stop_job(tid) + with JOBS_LOCK: + JOBS.pop(tid, None) + store.delete_task(tid) + return jsonify({"ok": True}) + + +# ---------------- API: 运行控制 ---------------- + +@app.route("/api/tasks//start", methods=["POST"]) +def api_start(tid): + task = store.get_task(tid) + if not task: + return jsonify({"error": "任务不存在"}), 404 + run, err = start_run(task) + if err: + return jsonify({"error": err}), 409 + return jsonify(run) + + +@app.route("/api/tasks//stop", methods=["POST"]) +def api_stop(tid): + if _stop_job(tid): + return jsonify({"ok": True, "msg": "已发送终止信号"}) + return jsonify({"ok": True, "msg": "任务未在运行"}) + + +@app.route("/api/tasks//pause", methods=["POST"]) +def api_pause(tid): + with JOBS_LOCK: + job = JOBS.get(tid) + if job and job.is_running(): + job.pause() + return jsonify({"ok": True}) + return jsonify({"error": "任务未在运行"}), 409 + + +@app.route("/api/tasks//resume", methods=["POST"]) +def api_resume(tid): + with JOBS_LOCK: + job = JOBS.get(tid) + if job and job.is_running(): + job.resume() + return jsonify({"ok": True}) + return jsonify({"error": "任务未在运行"}), 409 + + +# ---------------- API: 运行记录与文件 ---------------- + +@app.route("/api/runs/") +def api_run_detail(rid): + run = store.get_run(rid) + if not run: + return jsonify({"error": "运行记录不存在"}), 404 + return jsonify(run) + + +@app.route("/api/runs//logs") +def api_run_logs(rid): + run = store.get_run(rid) + if not run: + return jsonify({"error": "运行记录不存在"}), 404 + offset = int(request.args.get("offset", 0)) + logs = run.get("logs", []) + return jsonify({"logs": logs[offset:], "count": len(logs)}) + + +@app.route("/api/file") +def api_file(): + tid = request.args.get("task_id", "") + path = request.args.get("path", "") + task = store.get_task(tid) + if not task or not path: + return jsonify({"error": "参数错误"}), 400 + out_dir = os.path.realpath(resolve_out_dir(task)) + full = os.path.realpath(os.path.join(out_dir, path)) + if not full.startswith(out_dir + os.sep) and full != out_dir: + return jsonify({"error": "路径越界"}), 403 + if not os.path.isfile(full): + return jsonify({"error": "文件不存在"}), 404 + return send_file(full) + + +# ---------------- 启动 ---------------- + +scheduler = Scheduler(start_run) + + +def main(): + os.makedirs(os.path.join(HERE, "data"), exist_ok=True) + os.makedirs(os.path.join(HERE, "out"), exist_ok=True) + scheduler.start() + print(f"[universal-crawler] 启动完成, 管理界面: http://0.0.0.0:{PORT}/") + app.run(host="0.0.0.0", port=PORT, threaded=True, debug=False) + + +if __name__ == "__main__": + main() diff --git a/engine.py b/engine.py new file mode 100644 index 0000000..5946b17 --- /dev/null +++ b/engine.py @@ -0,0 +1,432 @@ +# -*- coding: utf-8 -*- +""" +通用爬虫引擎 (Playwright + stealth + 系统 Chrome) +- 批量模式: 逐条爬取网址列表 +- 自动模式: 从起始网址按匹配规则自动发现链接并 BFS 爬取 +- 支持: 随机延迟 / 重试 / 图片下载 / 暂停恢复 / 终止 / 配置热更新 / cookie 复用 +""" +import json +import os +import random +import re +import threading +import time +import urllib.parse +from datetime import datetime + +from playwright.sync_api import sync_playwright +from playwright_stealth import Stealth + +import notify +import store + +HERE = os.path.dirname(os.path.abspath(__file__)) +DATA_DIR = os.path.join(HERE, "data") +CHROME = "/usr/bin/google-chrome" +IMG_EXTS = (".jpg", ".jpeg", ".png", ".gif", ".webp", ".bmp") +DEFAULT_UA = ("Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 " + "(KHTML, like Gecko) Chrome/125.0.0.0 Safari/537.36") + +_CHALLENGE_MARKS = [ + "access denied", "403 forbidden", "just a moment", "attention required", + "captcha", "bot check", "cf-challenge", "verify you are human", + "checking your browser", "enable javascript and cookies", +] + + +def is_challenge_page(title, html): + low = html.lower() + t = (title or "").lower() + for mark in _CHALLENGE_MARKS: + if mark in t or mark in low: + return True + # 极小页面 + 无正文结构 -> 疑似验证壳 + if len(html) < 5000 and (" 500: + del logs[:len(logs) - 500] + self.persist(self.task["id"], self.run) + + def _persist(self): + self.persist(self.task["id"], self.run) + + def _resolve_out_dir(self): + cfg = self._cfg("out_dir") or "" + if cfg.strip(): + return cfg.strip() + return os.path.join(HERE, "out", self.task["id"]) + + def _open_browser(self): + p = sync_playwright().start() + browser = p.chromium.launch( + headless=True, + executable_path=CHROME, + args=["--disable-blink-features=AutomationControlled", + "--no-sandbox", "--disable-gpu", "--disable-dev-shm-usage"], + ) + ctx = browser.new_context( + user_agent=self._cfg("user_agent", DEFAULT_UA), + viewport={"width": 1920, "height": 1080}, + locale="en-US", + ) + cookie_file = os.path.join(DATA_DIR, f"cookies_{self.task['id']}.json") + if os.path.exists(cookie_file): + try: + ctx.add_cookies(json.load(open(cookie_file))) + self._log("info", "已复用上次会话 cookie") + except Exception: + pass + Stealth().apply_stealth_sync(ctx) + page = ctx.new_page() + return p, browser, ctx, page, cookie_file + + def _close_browser(self, p, browser, ctx, cookie_file): + try: + json.dump(ctx.cookies(), open(cookie_file, "w")) + except Exception: + pass + try: + browser.close() + except Exception: + pass + try: + p.stop() + except Exception: + pass + + def _wait_page_settle(self, page, timeout_s): + last_title, stable = "", 0 + start = time.time() + while time.time() - start < timeout_s: + if self._stop.is_set(): + raise RuntimeError("任务已终止") + self._wait_if_paused() + time.sleep(1) + try: + title = page.title() + html = page.content() + except Exception: + continue # 正在跳转 + if is_challenge_page(title, html): + stable = 0 + continue + if title == last_title: + stable += 1 + if stable >= 2 and len(html) > 1000: + return True, title, html + else: + stable = 0 + last_title = title + return True, page.title(), page.content() + + def _crawl_one(self, page, url, timeout_s): + page.goto(url, wait_until="domcontentloaded", timeout=timeout_s * 1000) + ok, title, html = self._wait_page_settle(page, timeout_s) + if not ok: + raise RuntimeError(f"页面加载超时({timeout_s}s)") + if is_challenge_page(title, html): + raise RuntimeError(f"仍被反爬拦截: title={title!r} size={len(html)}") + try: + text = page.inner_text("body") + except Exception: + text = "" + return title, html, text + + def _crawl_images(self, ctx, page, out_dir, base): + """下载页面图片, 返回 [{file, url, size}]""" + try: + urls = page.evaluate( + "() => Array.from(document.querySelectorAll('img'))" + ".map(i => i.currentSrc || i.src).filter(Boolean)" + ) + except Exception: + return [] + saved = [] + img_dir = os.path.join(out_dir, base + "_img") + for n, u in enumerate(urls, 1): + if self._stop.is_set(): + break + if not str(u).startswith("http"): + continue + path = urllib.parse.urlparse(u).path.lower() + if not path.endswith(IMG_EXTS): + continue + try: + resp = ctx.request.get(u, timeout=20000) + if resp.ok and resp.body(): + ext = os.path.splitext(path)[1] or ".jpg" + fname = f"img_{n:04d}{ext}" + os.makedirs(img_dir, exist_ok=True) + with open(os.path.join(img_dir, fname), "wb") as f: + f.write(resp.body()) + saved.append({"file": f"{base}_img/{fname}", "url": u, "size": len(resp.body())}) + except Exception: + continue + return saved + + def _retry_crawl(self, page, ctx, url, idx, out_dir): + """带重试的单页爬取, 返回结果 entry""" + timeout = int(self._cfg("timeout", 60)) + retries = int(self._cfg("retry_count", 2)) + retry_wait = float(self._cfg("retry_interval", 3)) + crawl_images = bool(self._cfg("crawl_images", False)) + + entry = {"url": url, "title": "", "status": "FAIL", "error": "", + "html_file": "", "txt_file": "", "images": []} + base = safe_name(url, idx) + for attempt in range(retries + 1): + if self._stop.is_set(): + entry["error"] = "任务已终止" + break + self._wait_if_paused() + try: + title, html, text = self._crawl_one(page, url, timeout) + html_path = os.path.join(out_dir, base + ".html") + txt_path = os.path.join(out_dir, base + ".txt") + with open(html_path, "w", encoding="utf-8") as f: + f.write(html) + with open(txt_path, "w", encoding="utf-8") as f: + f.write(text) + entry.update(title=title, status="OK", + html_file=base + ".html", txt_file=base + ".txt") + if crawl_images: + entry["images"] = self._crawl_images(ctx, page, out_dir, base) + self._log("info", f"OK {title[:50]!r} html={len(html)//1024}KB 图片={len(entry['images'])}") + break + except Exception as e: + entry["error"] = str(e) + self._log("warn", f"第{attempt + 1}次失败 {url}: {e}") + if attempt < retries: + self._wait_if_paused() + t0 = time.time() + while time.time() - t0 < retry_wait: + if self._stop.is_set(): + break + self._wait_if_paused() + time.sleep(0.3) + return entry + + def _delay(self): + dmin = float(self._cfg("delay_min", 2)) + dmax = float(self._cfg("delay_max", 5)) + total = random.uniform(max(0.1, dmin), max(dmin + 0.1, dmax)) + end = time.time() + total + while time.time() < end: + if self._stop.is_set(): + return + self._wait_if_paused() + time.sleep(min(0.5, end - time.time())) + + def _bump_stats(self, entry): + st = self.run.setdefault("stats", {"ok": 0, "fail": 0, "images": 0}) + if entry["status"] == "OK": + st["ok"] += 1 + else: + st["fail"] += 1 + st["images"] = st.get("images", 0) + len(entry.get("images", [])) + + # ---------------- 主流程 ---------------- + def _run_loop(self): + run, task = self.run, self.task + run["status"] = "running" + run["started_at"] = store.now_str() + self._persist() + try: + if task.get("mode") == "auto": + self._crawl_auto() + else: + self._crawl_list(task.get("urls", [])) + if self._stop.is_set(): + run["status"] = "stopped" + else: + run["status"] = "completed" + except Exception as e: + run["status"] = "failed" + self._log("error", f"任务异常终止: {e}") + run["finished_at"] = store.now_str() + self._persist() + if self._cfg("notify", False): + try: + ok, msg = notify.notify_email(task, run) + if ok: + self._log("info", "完成通知邮件已发送") + else: + self._log("error", f"邮件通知失败: {msg}") + except Exception as e: + self._log("error", f"邮件通知异常: {e}") + self._persist() + self._log("info", f"任务结束: {run['status']} 成功{run['stats'].get('ok', 0)} 失败{run['stats'].get('fail', 0)}") + self._persist() + + def _crawl_list(self, urls): + run = self.run + total = len(urls) + run["progress"]["total"] = total + out_dir = self._resolve_out_dir() + run["out_dir"] = out_dir + os.makedirs(out_dir, exist_ok=True) + self._persist() + + p, browser, ctx, page, cookie_file = self._open_browser() + try: + for i, url in enumerate(urls, 1): + if self._stop.is_set(): + self._log("info", "收到终止信号, 停止爬取") + break + self._wait_if_paused() + run["progress"]["current_url"] = url + run["progress"]["done"] = i - 1 + self._persist() + entry = self._retry_crawl(page, ctx, url, i, out_dir) + run["results"].append(entry) + self._bump_stats(entry) + run["progress"]["done"] = i + self._persist() + if entry["status"] == "OK": + self._delay() + finally: + self._close_browser(p, browser, ctx, cookie_file) + + def _discover_links(self, page): + """从当前页面提取符合规则的链接""" + auto = self.task.get("auto", {}) + include = [x.strip() for x in (auto.get("include") or []) if x.strip()] + exclude = [x.strip() for x in (auto.get("exclude") or []) if x.strip()] + use_regex = bool(auto.get("use_regex", False)) + same_domain = auto.get("same_domain", True) + seed_host = urllib.parse.urlparse(auto.get("seed_url", "")).hostname or "" + try: + hrefs = page.evaluate( + "() => Array.from(document.querySelectorAll('a[href]')).map(a => a.href)" + ) + except Exception: + return [] + out = [] + for h in hrefs: + if not str(h).startswith("http"): + continue + if same_domain and seed_host: + host = urllib.parse.urlparse(h).hostname or "" + if host != seed_host and not host.endswith("." + seed_host): + continue + if use_regex: + if include and not any(re.search(p, h) for p in include): + continue + if any(re.search(p, h) for p in exclude): + continue + else: + if include and not any(p.lower() in h.lower() for p in include): + continue + if any(p.lower() in h.lower() for p in exclude): + continue + out.append(h) + return out + + def _crawl_auto(self): + run = self.run + auto = self.task.get("auto", {}) + seed = auto.get("seed_url", "") + max_pages = int(auto.get("max_pages", 50) or 50) + max_depth = int(auto.get("max_depth", 2) or 2) + run["progress"]["total"] = max_pages + out_dir = self._resolve_out_dir() + run["out_dir"] = out_dir + os.makedirs(out_dir, exist_ok=True) + self._persist() + + p, browser, ctx, page, cookie_file = self._open_browser() + queue = [(seed, 0)] + visited = set() + queued = set([seed]) + idx = 0 + try: + while queue and not self._stop.is_set(): + self._wait_if_paused() + url, depth = queue.pop(0) + if url in visited: + continue + if len(visited) >= max_pages: + break + visited.add(url) + idx += 1 + run["progress"]["current_url"] = url + run["progress"]["done"] = len(visited) + self._persist() + entry = self._retry_crawl(page, ctx, url, idx, out_dir) + run["results"].append(entry) + self._bump_stats(entry) + self._persist() + if entry["status"] == "OK" and depth < max_depth: + for link in self._discover_links(page): + if link not in visited and link not in queued: + queued.add(link) + queue.append((link, depth + 1)) + if entry["status"] == "OK": + self._delay() + finally: + self._close_browser(p, browser, ctx, cookie_file) + run["progress"]["total"] = len(visited) + self._persist() diff --git a/legacy/crawl.py b/legacy/crawl.py new file mode 100644 index 0000000..9ee7ac0 --- /dev/null +++ b/legacy/crawl.py @@ -0,0 +1,168 @@ +# -*- coding: utf-8 -*- +""" +批量网页爬取器 (Playwright + stealth + 系统Chrome) +用法: + python crawl.py urls.txt [--out 输出目录] [--delay 2,5] [--timeout 60] + +urls.txt: 每行一个网址, # 开头为注释, 空行忽略 +输出: + out/0001_<域名>_<时间戳>.html 完整网页HTML + out/0001_<域名>_<时间戳>.txt 提取的纯文本正文 + out/results.csv 汇总表(序号/网址/标题/状态/文件) + out/cookies.json 会话cookie(自动复用, 减少重复验证) +""" +import argparse, csv, json, os, random, re, sys, time, urllib.parse +from datetime import datetime +from playwright.sync_api import sync_playwright +from playwright_stealth import Stealth + +HERE = os.path.dirname(os.path.abspath(__file__)) +COOKIE_FILE = os.path.join(HERE, "cookies.json") + +def load_urls(path): + urls = [] + with open(path, encoding="utf-8") as f: + for line in f: + line = line.strip() + if not line or line.startswith("#"): + continue + if not line.startswith("http"): + line = "https://" + line + urls.append(line) + return urls + +def safe_name(url, idx): + host = urllib.parse.urlparse(url).netloc.replace("www.", "").replace(".", "_") + ts = datetime.now().strftime("%Y%m%d_%H%M%S") + return f"{idx:04d}_{host}_{ts}" + +def is_challenge_page(title, html): + """判断是否仍在反爬验证页 (通用启发式, 可按需调整)""" + low = html.lower() + if "access denied" in title.lower() or "403" in title: + return True + if "bot check" in low or "captcha" in low or "cf-challenge" in low: + return True + # 页面过小且标题是站名本身 -> 多半是验证壳 + if len(html) < 8000 and "techpowerup" in low: + return True + return False + +def wait_page_settle(page, timeout_s=60): + """等页面稳定: 挑战自动跳转结束 + 内容可读""" + last_title, stable = "", 0 + start = time.time() + while time.time() - start < timeout_s: + time.sleep(1) + try: + title = page.title() + html = page.content() + except Exception: + continue # 正在跳转 + if is_challenge_page(title, html): + stable = 0 + continue + if title == last_title: + stable += 1 + if stable >= 2 and len(html) > 1000: + return True, title, html + else: + stable = 0 + last_title = title + try: + return True, page.title(), page.content() + except Exception: + return False, "", "" + +def crawl_one(page, url, timeout_s): + page.goto(url, wait_until="domcontentloaded", timeout=timeout_s * 1000) + ok, title, html = wait_page_settle(page, timeout_s) + if not ok: + raise RuntimeError(f"页面加载超时({timeout_s}s)") + if is_challenge_page(title, html): + raise RuntimeError(f"仍被反爬拦截: title={title!r} size={len(html)}") + text = page.inner_text("body") + return title, html, text + +def main(): + ap = argparse.ArgumentParser() + ap.add_argument("urls_file") + ap.add_argument("--out", default=os.path.join(HERE, "out")) + ap.add_argument("--delay", default="2,5", help="每次请求间随机延迟秒数, 如 2,5") + ap.add_argument("--timeout", type=int, default=60, help="单页最长等待秒数") + ap.add_argument("--retry", type=int, default=2, help="失败重试次数") + args = ap.parse_args() + + urls = load_urls(args.urls_file) + if not urls: + print("urls.txt 里没有有效网址"); return + os.makedirs(args.out, exist_ok=True) + dmin, dmax = map(float, args.delay.split(",")) + + results = [] + with sync_playwright() as p: + browser = p.chromium.launch( + headless=True, + executable_path="/usr/bin/google-chrome", + args=["--disable-blink-features=AutomationControlled", + "--no-sandbox", "--disable-gpu"], + ) + ctx = browser.new_context( + user_agent="Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 " + "(KHTML, like Gecko) Chrome/125.0.0.0 Safari/537.36", + viewport={"width": 1920, "height": 1080}, + locale="en-US", + ) + if os.path.exists(COOKIE_FILE): + try: + ctx.add_cookies(json.load(open(COOKIE_FILE))) + print("[i] 已复用上次会话 cookie") + except Exception as e: + print(f"[i] cookie 复用失败(忽略): {e}") + Stealth().apply_stealth_sync(ctx) + page = ctx.new_page() + + for i, url in enumerate(urls, 1): + print(f"[{i}/{len(urls)}] {url}", flush=True) + base = safe_name(url, i) + title, status = "", "" + for attempt in range(args.retry + 1): + try: + title, html, text = crawl_one(page, url, args.timeout) + status = "OK" + html_path = os.path.join(args.out, base + ".html") + txt_path = os.path.join(args.out, base + ".txt") + with open(html_path, "w", encoding="utf-8") as f: f.write(html) + with open(txt_path, "w", encoding="utf-8") as f: f.write(text) + print(f" -> OK {title[:60]!r} html={len(html)//1024}KB txt={len(text)//1024}KB") + break + except Exception as e: + status = f"FAIL: {e}" + print(f" -> 第{attempt+1}次失败: {e}") + time.sleep(3) + results.append({"no": i, "url": url, "title": title, "status": status, + "file": base + ".html"}) + + # 保存 cookie (验证通过后), 供下次复用 + try: + json.dump(ctx.cookies(), open(COOKIE_FILE, "w")) + except Exception: + pass + time.sleep(random.uniform(dmin, dmax)) # 礼貌延迟 + + browser.close() + + # 汇总 + csv_path = os.path.join(args.out, "results.csv") + with open(csv_path, "w", newline="", encoding="utf-8-sig") as f: + w = csv.writer(f) + w.writerow(["序号", "网址", "标题", "状态", "文件"]) + for r in results: + w.writerow([r["no"], r["url"], r["title"], r["status"], r["file"]]) + ok = sum(1 for r in results if r["status"] == "OK") + print(f"\n完成: {ok}/{len(results)} 成功") + print(f"HTML/文本: {args.out}/") + print(f"汇总表: {csv_path}") + +if __name__ == "__main__": + main() diff --git a/legacy/test_one.py b/legacy/test_one.py new file mode 100644 index 0000000..b4ba88e --- /dev/null +++ b/legacy/test_one.py @@ -0,0 +1,48 @@ +# -*- coding: utf-8 -*- +"""单页测试: Playwright + stealth + 系统Chrome 抓 TPU""" +import sys, time +from playwright.sync_api import sync_playwright +from playwright_stealth import Stealth + +URL = sys.argv[1] if len(sys.argv) > 1 else "https://www.techpowerup.com/" + +with sync_playwright() as p: + browser = p.chromium.launch( + headless=True, + executable_path="/usr/bin/google-chrome", + args=["--disable-blink-features=AutomationControlled", "--no-sandbox", "--disable-gpu"], + ) + ctx = browser.new_context( + user_agent="Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/125.0.0.0 Safari/537.36", + viewport={"width": 1920, "height": 1080}, + locale="en-US", + ) + Stealth().apply_stealth_sync(ctx) # 隐藏 webdriver/headless 特征 + page = ctx.new_page() + try: + page.goto(URL, wait_until="domcontentloaded", timeout=60000) + # 等 JS 挑战自动通过: 最多等 40s, 轮询检查标题/内容特征 + for i in range(40): + time.sleep(1) + try: + title = page.title() + html = page.content() + except Exception: + continue # 页面正在跳转, 稍后再试 + if "Access Denied" in title or "403" in title: + print(f"[{i}s] 403 Access Denied"); break + if "bot check" in html.lower() and len(html) < 5000: + continue + if len(html) > 30000 or "news" in html.lower(): + print(f"[{i}s] 通过! title={title!r} html_size={len(html)}") + break + else: + print("超时未通过, 当前 title:", page.title()) + # 打印正文摘要作为证据 + text = page.inner_text("body")[:800] + print("--- body 前800字符 ---") + print(text) + except Exception as e: + print("ERROR:", e) + finally: + browser.close() diff --git a/legacy/urls.txt b/legacy/urls.txt new file mode 100644 index 0000000..8e6507e --- /dev/null +++ b/legacy/urls.txt @@ -0,0 +1,4 @@ +# 测试用网址列表 (TPU 首页 + 一篇评测 + GPU数据库页) +https://www.techpowerup.com/ +https://www.techpowerup.com/review/xmg-apex-17-m25/ +https://www.techpowerup.com/gpu-specs/ diff --git a/notify.py b/notify.py new file mode 100644 index 0000000..1824931 --- /dev/null +++ b/notify.py @@ -0,0 +1,43 @@ +# -*- coding: utf-8 -*- +"""邮件通知 (复用 send_email.py 技能脚本)""" +import os +import subprocess +import sys + +SEND_EMAIL = os.path.join( + os.path.dirname(os.path.dirname(os.path.abspath(__file__))), + "skills", "send_email.py", +) + + +def notify_email(task, run): + """发送完成通知邮件, 返回 (bool, msg)""" + if not os.path.exists(SEND_EMAIL): + return False, "send_email.py 不存在" + cfg = task.get("config", {}) + to = cfg.get("notify_email") or "wlq@tphai.com" + results = run.get("results", []) + ok = sum(1 for r in results if r.get("status") == "OK") + fail = len(results) - ok + lines = [ + f"项目: {task['name']}", + f"模式: {task['mode']}", + f"状态: {run.get('status')}", + f"成功: {ok} 失败: {fail} 图片: {run.get('stats', {}).get('images', 0)}", + f"输出目录: {run.get('out_dir', '')}", + "", + "明细:", + ] + for r in results[:50]: + lines.append(f" [{r.get('status')}] {r.get('title', '')[:40]} {r.get('url', '')}") + if len(results) > 50: + lines.append(f" ... 共 {len(results)} 条") + body = "\n".join(lines) + try: + r = subprocess.run( + [sys.executable, SEND_EMAIL, f"[爬虫完成] {task['name']}", body, "--to", to], + timeout=60, capture_output=True, + ) + return r.returncode == 0, (r.stderr or b"").decode(errors="ignore")[:200] + except Exception as e: + return False, str(e) diff --git a/requirements.txt b/requirements.txt new file mode 100644 index 0000000..79264a8 --- /dev/null +++ b/requirements.txt @@ -0,0 +1,3 @@ +flask>=3.0 +playwright>=1.40 +playwright-stealth>=1.0 diff --git a/scheduler.py b/scheduler.py new file mode 100644 index 0000000..438b9b2 --- /dev/null +++ b/scheduler.py @@ -0,0 +1,120 @@ +# -*- coding: utf-8 -*- +"""定时调度器: 支持间隔调度与 5 段 cron 表达式""" +import threading +from datetime import datetime, timedelta + +import store + + +# ---------------- cron 工具 ---------------- + +def _field(f, lo, hi): + out = set() + for part in str(f).split(","): + part = part.strip() + if not part: + continue + if "/" in part: + rng, step = part.split("/") + step = int(step) + if rng == "*": + out.update(range(lo, hi + 1, step)) + elif "-" in rng: + a, b = map(int, rng.split("-")) + out.update(range(a, b + 1, step)) + else: + out.update(range(int(rng), hi + 1, step)) + elif "-" in part: + a, b = map(int, part.split("-")) + out.update(range(a, b + 1)) + elif part == "*": + out.update(range(lo, hi + 1)) + else: + out.add(int(part)) + return out + + +def cron_next(expr, base=None): + """计算 cron 表达式(分 时 日 月 周, 周:0/7=周日)的下一次执行时间""" + base = base or datetime.now() + parts = str(expr).strip().split() + if len(parts) != 5: + raise ValueError("cron 表达式必须为 5 段: 分 时 日 月 周") + mins = _field(parts[0], 0, 59) + hours = _field(parts[1], 0, 23) + doms = _field(parts[2], 1, 31) + mons = _field(parts[3], 1, 12) + dows = _field(parts[4], 0, 7) + dows = {d % 7 for d in dows} # cron 0/7=周日 -> python 0=周一..6=周日 + t = base.replace(second=0, microsecond=0) + timedelta(minutes=1) + for _ in range(366 * 24 * 60): # 一年窗口 + if (t.minute in mins and t.hour in hours and t.day in doms + and t.month in mons and t.weekday() in dows): + return t + t += timedelta(minutes=1) + return None + + +def interval_delta(sch): + unit = sch.get("interval_unit", "hours") + val = int(sch.get("interval_value", 24) or 24) + if unit == "minutes": + return timedelta(minutes=val) + if unit == "days": + return timedelta(days=val) + return timedelta(hours=val) + + +# ---------------- 调度器线程 ---------------- + +class Scheduler(threading.Thread): + """每 30 秒检查一次所有定时任务, 到点则触发一次爬取""" + + def __init__(self, spawn): + """ + spawn: callable(task) -> 启动一次运行 (由 app 提供) + """ + super().__init__(daemon=True, name="scheduler") + self._spawn = spawn + self._stop = threading.Event() + + def stop(self): + self._stop.set() + + def run(self): + while not self._stop.wait(30): + try: + self._tick() + except Exception: + pass + + def _tick(self): + now = datetime.now() + for task in store.load_tasks(): + sch = task.get("schedule") or {} + if task.get("mode") != "scheduled" or not sch.get("enabled"): + continue + nxt = sch.get("next_run") + if not nxt: + continue + try: + due = datetime.strptime(nxt, "%Y-%m-%d %H:%M:%S") + except Exception: + continue + if due > now: + continue + runs = store.get_runs(task["id"]) + if runs and runs[-1].get("status") == "running": + continue # 上次还没跑完, 跳过本次 + self._spawn(task) + try: + if sch.get("type") == "cron": + nn = cron_next(sch.get("cron", "0 * * * *"), now) + else: + nn = now + interval_delta(sch) + sch["next_run"] = nn.strftime("%Y-%m-%d %H:%M:%S") if nn else None + sch["last_run"] = now.strftime("%Y-%m-%d %H:%M:%S") + sch["runs_count"] = sch.get("runs_count", 0) + 1 + store.upsert_task(task) + except Exception: + pass diff --git a/start.sh b/start.sh new file mode 100755 index 0000000..8e58610 --- /dev/null +++ b/start.sh @@ -0,0 +1,54 @@ +#!/bin/bash +# 通用爬虫系统 - 启动/停止脚本 (PID 文件管理) +cd "$(dirname "$0")" +PY=/home/hz1/miniconda3/envs/openclaw/bin/python +PORT=${CRAWLER_PORT:-16062} +PIDFILE=logs/app.pid + +get_pid() { + if [ -f "$PIDFILE" ]; then + cat "$PIDFILE" 2>/dev/null + fi +} + +is_running() { + local pid=$(get_pid) + [ -n "$pid" ] && kill -0 "$pid" 2>/dev/null +} + +case "$1" in + stop) + if is_running; then + kill $(get_pid) && rm -f "$PIDFILE" && echo "已停止" + else + rm -f "$PIDFILE"; echo "未在运行" + fi + ;; + restart) + if is_running; then kill $(get_pid); sleep 1; fi + rm -f "$PIDFILE" + mkdir -p logs + nohup $PY app.py > logs/app.log 2>&1 & + echo $! > "$PIDFILE" + sleep 1 + echo "已重启, PID $(get_pid) (端口 $PORT)" + ;; + status) + if is_running; then + echo "运行中, PID $(get_pid) (端口 $PORT)" + else + echo "未运行" + fi + ;; + *) + if is_running; then + echo "已在运行, PID $(get_pid) (端口 $PORT)"; exit 0 + fi + mkdir -p logs + nohup $PY app.py > logs/app.log 2>&1 & + echo $! > "$PIDFILE" + sleep 1 + echo "已启动, PID $(get_pid) (端口 $PORT)" + echo "管理界面: http://$(hostname -I | awk '{print $1}'):$PORT/" + ;; +esac diff --git a/static/app.js b/static/app.js new file mode 100644 index 0000000..6a4a1aa --- /dev/null +++ b/static/app.js @@ -0,0 +1,497 @@ +/* 通用爬虫系统 - 前端逻辑 */ +"use strict"; + +const $ = (id) => document.getElementById(id); +const MODE_LABEL = { batch: "批量", scheduled: "定时", auto: "自动" }; + +const state = { + tasks: [], + editTask: null, // 正在编辑的任务 + mode: "batch", // 当前表单模式 + detail: { task: null, runId: null, logOffset: 0, timer: null }, + logTimer: null, +}; + +/* ---------------- 工具 ---------------- */ +function toast(msg, isErr) { + const t = document.createElement("div"); + t.className = "toast" + (isErr ? " error" : ""); + t.textContent = msg; + document.body.appendChild(t); + setTimeout(() => t.remove(), 2600); +} + +async function api(url, opts) { + const res = await fetch(url, opts); + let data = null; + try { data = await res.json(); } catch (e) { /* ignore */ } + if (!res.ok) throw new Error((data && data.error) || `HTTP ${res.status}`); + return data; +} + +const esc = (s) => String(s == null ? "" : s).replace(/[&<>"']/g, + (c) => ({ "&": "&", "<": "<", ">": ">", '"': """, "'": "'" }[c])); + +function fmtTime(iso) { return iso ? String(iso).replace("T", " ").slice(0, 19) : "—"; } + +/* ---------------- 任务列表 ---------------- */ +async function loadTasks() { + try { + state.tasks = await api("/api/tasks"); + renderTasks(); + const running = state.tasks.filter((t) => t.running).length; + const pill = $("statusPill"); + pill.textContent = `运行中任务: ${running}`; + pill.classList.toggle("active", running > 0); + const st = await api("/api/status"); + $("version").textContent = `v${st.version}`; + } catch (e) { + toast("加载任务失败: " + e.message, true); + } +} + +function taskStatusBadge(t) { + if (t.running) return '● 运行中'; + const r = t.latest_run; + if (r && r.status === "paused") return '⏸ 已暂停'; + if (t.mode === "scheduled" && t.schedule && t.schedule.enabled) { + const next = t.schedule.next_run ? ` · 下次 ${t.schedule.next_run.slice(11)}` : ""; + return `⏰ 等待调度${esc(next)}`; + } + if (r) return `${({ completed: "✅ 已完成", failed: "❌ 失败", stopped: "⏹ 已停止" })[r.status] || esc(r.status)}`; + return '未运行'; +} + +function renderTasks() { + const box = $("taskList"); + $("emptyState").classList.toggle("hidden", state.tasks.length > 0); + box.innerHTML = state.tasks.map(taskCard).join(""); +} + +function taskCard(t) { + const r = t.latest_run; + const cfg = t.config || {}; + const prog = r && (r.status === "running" || r.status === "paused") ? ` +
+
+ ${r.progress.done || 0} / ${r.progress.total || "?"} + ${esc(r.progress.current_url || "")} +
` : ""; + const stats = r ? ` +
+ ✅ ${r.stats.ok || 0} + ❌ ${r.stats.fail || 0} + 🖼️ ${r.stats.images || 0} +
` : ""; + const urlCount = (t.urls || []).length; + const line2 = t.mode === "auto" + ? `起始: ${esc((t.auto && t.auto.seed_url) || "")}` + : `网址: ${urlCount} 个`; + const line3 = t.mode === "scheduled" + ? `调度: ${t.schedule.type === "cron" ? esc(t.schedule.cron) : "每 " + (t.schedule.interval_value || "-") + " " + ({ minutes: "分钟", hours: "小时", days: "天" }[t.schedule.interval_unit] || "")}` + : ""; + const outDir = cfg.out_dir || "out/" + t.id; + + let acts = ""; + if (t.running) { + acts = ` + + `; + } else { + const canStart = t.mode !== "scheduled" || (t.urls || []).length > 0; + acts = ``; + } + acts += ` + + + `; + + return ` +
+
+ ${esc(t.name)} + ${taskStatusBadge(t)} +
+
+ ${MODE_LABEL[t.mode]} + ${line2} +
+ ${line3 ? `
${line3}
` : ""} +
输出: ${esc(outDir)}
+ ${prog} + ${stats} +
${acts}
+
`; +} + +/* ---------------- 任务操作 ---------------- */ +async function actTask(tid, act) { + try { + const r = await api(`/api/tasks/${tid}/${act}`, { method: "POST" }); + if (r.msg) toast(r.msg); + loadTasks(); + if (state.detail.task && state.detail.task.id === tid) openDetail(tid); + } catch (e) { toast(e.message, true); } +} + +async function delTask(tid) { + const t = state.tasks.find((x) => x.id === tid); + if (!confirm(`确定删除任务「${t ? t.name : tid}」?\n运行记录与任务配置将被删除(已爬取的文件保留在磁盘上)。`)) return; + try { + await api(`/api/tasks/${tid}`, { method: "DELETE" }); + toast("已删除"); + loadTasks(); + } catch (e) { toast(e.message, true); } +} + +/* ---------------- 新建 / 编辑表单 ---------------- */ +function switchMode(mode, lock) { + state.mode = mode; + document.querySelectorAll("#modeTabs .tab").forEach((b) => { + b.classList.toggle("active", b.dataset.mode === mode); + b.disabled = !!lock; + }); + $("urlsField").classList.toggle("hidden", mode !== "batch" && mode !== "scheduled"); + $("scheduleBox").classList.toggle("hidden", mode !== "scheduled"); + $("autoBox").classList.toggle("hidden", mode !== "auto"); +} + +function openCreate() { + state.editTask = null; + $("modalTitle").textContent = "新建爬取任务"; + $("taskForm").reset(); + $("taskForm").elements["delay_min"].value = 2; + $("taskForm").elements["delay_max"].value = 5; + $("taskForm").elements["timeout"].value = 60; + $("taskForm").elements["retry_count"].value = 2; + $("taskForm").elements["retry_interval"].value = 3; + $("taskForm").elements["notify_email"].value = "wlq@tphai.com"; + $("formHint").textContent = ""; + switchMode("batch", false); + showModal("taskModal"); +} + +async function openEdit(tid) { + const t = state.tasks.find((x) => x.id === tid) || await api(`/api/tasks/${tid}`); + state.editTask = t; + $("modalTitle").textContent = `编辑任务「${t.name}」`; + const f = $("taskForm"); + f.reset(); + const cfg = t.config || {}; + f.elements["name"].value = t.name || ""; + f.elements["out_dir"].value = cfg.out_dir || ""; + f.elements["delay_min"].value = cfg.delay_min ?? 2; + f.elements["delay_max"].value = cfg.delay_max ?? 5; + f.elements["timeout"].value = cfg.timeout ?? 60; + f.elements["retry_count"].value = cfg.retry_count ?? 2; + f.elements["retry_interval"].value = cfg.retry_interval ?? 3; + f.elements["crawl_images"].checked = !!cfg.crawl_images; + f.elements["notify"].checked = !!cfg.notify; + f.elements["notify_email"].value = cfg.notify_email || "wlq@tphai.com"; + f.elements["urls"].value = (t.urls || []).join("\n"); + if (t.mode === "scheduled" && t.schedule) { + f.elements["schedule_type"].value = t.schedule.type || "interval"; + f.elements["schedule_enabled"].checked = t.schedule.enabled !== false; + f.elements["interval_value"].value = t.schedule.interval_value ?? 24; + f.elements["interval_unit"].value = t.schedule.interval_unit || "hours"; + f.elements["cron"].value = t.schedule.cron || ""; + } + if (t.mode === "auto" && t.auto) { + f.elements["seed_url"].value = t.auto.seed_url || ""; + f.elements["include"].value = (t.auto.include || []).join("\n"); + f.elements["exclude"].value = (t.auto.exclude || []).join("\n"); + f.elements["same_domain"].checked = t.auto.same_domain !== false; + f.elements["use_regex"].checked = !!t.auto.use_regex; + f.elements["max_pages"].value = t.auto.max_pages ?? 50; + f.elements["max_depth"].value = t.auto.max_depth ?? 2; + } + syncScheduleUI(); + $("formHint").textContent = t.running + ? "⚠️ 任务运行中:参数修改将热更新(网址/规则改动下次运行生效)" + : ""; + switchMode(t.mode, true); + showModal("taskModal"); +} + +function syncScheduleUI() { + const f = $("taskForm"); + const t = f.elements["schedule_type"].value; + $("intervalBox").classList.toggle("hidden", t !== "interval"); + $("cronBox").classList.toggle("hidden", t !== "cron"); +} + +function splitLines(v) { return String(v || "").split("\n").map((s) => s.trim()).filter(Boolean); } + +async function submitForm(e) { + e.preventDefault(); + const f = $("taskForm"); + const name = f.elements["name"].value.trim(); + if (!name) { toast("请填写项目名称", true); return; } + const mode = state.mode; + const config = { + out_dir: f.elements["out_dir"].value.trim(), + delay_min: parseFloat(f.elements["delay_min"].value) || 2, + delay_max: parseFloat(f.elements["delay_max"].value) || 5, + timeout: parseInt(f.elements["timeout"].value) || 60, + retry_count: parseInt(f.elements["retry_count"].value) || 0, + retry_interval: parseFloat(f.elements["retry_interval"].value) || 3, + crawl_images: f.elements["crawl_images"].checked, + notify: f.elements["notify"].checked, + notify_email: f.elements["notify_email"].value.trim() || "wlq@tphai.com", + }; + const urls = splitLines(f.elements["urls"].value); + const payload = { name, mode, config, urls }; + if (mode === "auto") { + payload.auto = { + seed_url: f.elements["seed_url"].value.trim(), + include: splitLines(f.elements["include"].value), + exclude: splitLines(f.elements["exclude"].value), + same_domain: f.elements["same_domain"].checked, + use_regex: f.elements["use_regex"].checked, + max_pages: parseInt(f.elements["max_pages"].value) || 50, + max_depth: parseInt(f.elements["max_depth"].value) || 2, + }; + } + if (mode === "scheduled") { + payload.schedule = { + type: f.elements["schedule_type"].value, + enabled: f.elements["schedule_enabled"].checked, + interval_unit: f.elements["interval_unit"].value, + interval_value: parseInt(f.elements["interval_value"].value) || 1, + cron: f.elements["cron"].value.trim(), + }; + } + try { + if (state.editTask) { + await api(`/api/tasks/${state.editTask.id}`, { + method: "PUT", headers: { "Content-Type": "application/json" }, + body: JSON.stringify(payload), + }); + toast("任务已更新"); + } else { + const t = await api("/api/tasks", { + method: "POST", headers: { "Content-Type": "application/json" }, + body: JSON.stringify(payload), + }); + toast(`任务「${t.name}」已创建`); + } + hideModal("taskModal"); + loadTasks(); + } catch (err) { toast(err.message, true); } +} + +/* ---------------- 详情 ---------------- */ +async function openDetail(tid) { + try { + const t = await api(`/api/tasks/${tid}`); + state.detail.task = t; + const runs = t.runs || []; + state.detail.runId = runs.find((r) => r.status === "running") + ? runs[0].id : (runs[0] ? runs[0].id : null); + state.detail.logOffset = 0; + $("detailTitle").textContent = `任务详情 · ${t.name}`; + renderDetail(); + showModal("detailModal"); + startLogPoll(); + } catch (e) { toast(e.message, true); } +} + +function renderDetail() { + const t = state.detail.task; + const runs = t.runs || []; + const run = runs.find((r) => r.id === state.detail.runId) || runs[0] || null; + state.detail.runId = run ? run.id : null; + const cfg = t.config || {}; + + const meta = ` +
+ 模式: ${MODE_LABEL[t.mode] || t.mode} + 任务ID: ${t.id} + 创建: ${fmtTime(t.created_at)} + 输出: ${esc(cfg.out_dir || "out/" + t.id)} + ${t.mode === "auto" ? `起始: ${esc((t.auto && t.auto.seed_url) || "")}` : ""} +
`; + + let sched = ""; + if (t.mode === "scheduled" && t.schedule) { + const s = t.schedule; + sched = `
+ 调度: ${s.type === "cron" ? esc(s.cron) : "每 " + (s.interval_value || "-") + " " + ({ minutes: "分钟", hours: "小时", days: "天" }[s.interval_unit] || "")} + 启用: ${s.enabled ? "是" : "否"} + 下次: ${esc(s.next_run || "—")} + 上次: ${esc(s.last_run || "—")} + 已运行: ${s.runs_count || 0} +
`; + } + + const acts = ` +
+ ${t.running ? ` + + ` + : ``} + +
`; + + const chips = runs.map((r) => ` + + ${fmtTime(r.started_at).slice(5)} · ${r.status} · ${r.progress.done || 0}/${r.progress.total || "?"} + `).join(""); + + const body = run ? renderRunPanel(t, run) : '
暂无运行记录,点击「立即执行」开始第一次爬取。
'; + + $("detailBody").innerHTML = ` + ${meta} + ${sched} + ${acts} +
${chips || '暂无运行记录'}
+ ${body}`; +} + +function renderRunPanel(t, run) { + const s = run.stats || { ok: 0, fail: 0, images: 0 }; + const prog = run.status === "running" || run.status === "paused" ? ` +
` : ""; + const results = run.results || []; + const rows = results.map((r, i) => { + const statusCls = r.status === "OK" ? "t-ok" : "t-fail"; + const html = r.html_file ? `HTML` : "—"; + const txt = r.txt_file ? `TXT` : "—"; + const imgs = (r.images || []).length + ? `🖼️ ${(r.images || []).length}` : "—"; + return ` + ${i + 1} + ${r.status} + ${esc(r.title)} + ${esc(r.url)} + ${html}${txt}${imgs} + ${esc((r.error || "").slice(0, 60))} + `; + }).join(""); + const thumbs = results.flatMap((r) => r.images || []).slice(0, 60); + const thumbHtml = thumbs.length ? ` +
+ ${thumbs.map((im) => ``).join("")} +
` : ""; + + return ` +
+
+ 运行ID: ${run.id} + 状态: ${run.status} + 开始: ${fmtTime(run.started_at)} + 结束: ${fmtTime(run.finished_at)} + ✅ ${s.ok} + ❌ ${s.fail} + 🖼️ ${s.images} +
+ ${prog} +
当前: ${esc(run.progress.current_url || "")}
+ ${results.length ? ` +
+ + + ${rows} +
#状态标题网址HTMLTXT图片错误
+
` : '
暂无结果
'} + ${thumbHtml} +
+
`; +} + +function scrollThumbs() { + const box = $("thumbsBox"); + if (box) box.scrollIntoView({ behavior: "smooth", block: "center" }); +} + +async function selectRun(rid) { + state.detail.runId = rid; + state.detail.logOffset = 0; + renderDetail(); + startLogPoll(); +} + +async function startLogPoll() { + stopLogPoll(); + const t = state.detail.task; + if (!t || !state.detail.runId) return; + const rid = state.detail.runId; + const tick = async () => { + try { + const d = await api(`/api/runs/${rid}/logs?offset=${state.detail.logOffset}`); + if (d.logs && d.logs.length) { + const box = $("logBox"); + if (box) { + const atBottom = box.scrollHeight - box.scrollTop - box.clientHeight < 40; + d.logs.forEach((l) => { + const div = document.createElement("div"); + div.className = l.level; + div.textContent = `[${l.ts}] ${l.msg}`; + box.appendChild(div); + }); + if (atBottom) box.scrollTop = box.scrollHeight; + state.detail.logOffset += d.logs.length; + } + } + } catch (e) { /* ignore */ } + loadTasks(); // 后台保持列表刷新 + }; + tick(); + state.logTimer = setInterval(tick, 2500); +} + +function stopLogPoll() { + if (state.logTimer) { clearInterval(state.logTimer); state.logTimer = null; } +} + +/* ---------------- 文件预览 ---------------- */ +function previewFile(tid, path, title) { + $("previewTitle").textContent = title || "预览"; + $("previewFrame").src = `/api/file?task_id=${tid}&path=${encodeURIComponent(path)}`; + showModal("previewModal"); +} + +/* ---------------- 弹窗控制 ---------------- */ +function showModal(id) { $(id).classList.remove("hidden"); } +function hideModal(id) { $(id).classList.add("hidden"); } + +/* ---------------- 事件绑定 ---------------- */ +$("btnNew").onclick = openCreate; +$("btnNew2").onclick = openCreate; +$("btnRefresh").onclick = loadTasks; + +document.querySelectorAll("#modeTabs .tab").forEach((b) => { + b.onclick = () => { switchMode(b.dataset.mode, false); }; +}); +$("taskForm").onsubmit = submitForm; +$("taskForm").elements["schedule_type"].onchange = syncScheduleUI; + +document.querySelectorAll("[data-close]").forEach((b) => b.onclick = () => hideModal("taskModal")); +document.querySelectorAll("[data-close-detail]").forEach((b) => b.onclick = () => { stopLogPoll(); hideModal("detailModal"); }); +document.querySelectorAll("[data-close-preview]").forEach((b) => b.onclick = () => { $("previewFrame").src = "about:blank"; hideModal("previewModal"); }); + +document.querySelectorAll(".modal-overlay").forEach((ov) => { + ov.addEventListener("mousedown", (e) => { + if (e.target === ov) { + if (ov.id === "detailModal") stopLogPoll(); + if (ov.id === "previewModal") $("previewFrame").src = "about:blank"; + ov.classList.add("hidden"); + } + }); +}); +document.addEventListener("keydown", (e) => { + if (e.key === "Escape") { + stopLogPoll(); + $("previewFrame").src = "about:blank"; + ["taskModal", "detailModal", "previewModal"].forEach((id) => $(id).classList.add("hidden")); + } +}); + +/* ---------------- 启动 ---------------- */ +loadTasks(); +setInterval(() => { + const detailOpen = !$("detailModal").classList.contains("hidden"); + if (!detailOpen) loadTasks(); +}, 4000); diff --git a/static/index.html b/static/index.html new file mode 100644 index 0000000..b8d2463 --- /dev/null +++ b/static/index.html @@ -0,0 +1,143 @@ + + + + + +通用爬虫系统 + + + +
+ +
+ 运行中任务: 0 + + +
+
+ +
+
+ +
+ + + + + + + + + + + + + diff --git a/static/style.css b/static/style.css new file mode 100644 index 0000000..8895666 --- /dev/null +++ b/static/style.css @@ -0,0 +1,181 @@ +:root { + --bg: #0f1420; + --panel: #171e2e; + --panel2: #1e2740; + --border: #2a3550; + --text: #dbe4f5; + --muted: #8b97b5; + --accent: #4f8cff; + --green: #37d39a; + --red: #ff5d6c; + --yellow: #ffc857; + --purple: #a78bfa; +} +* { box-sizing: border-box; margin: 0; padding: 0; } +body { + background: var(--bg); color: var(--text); + font-family: -apple-system, "PingFang SC", "Microsoft YaHei", "Segoe UI", sans-serif; + font-size: 14px; min-height: 100vh; +} +header { + display: flex; align-items: center; justify-content: space-between; + padding: 14px 24px; background: var(--panel); + border-bottom: 1px solid var(--border); position: sticky; top: 0; z-index: 10; +} +.logo { font-size: 18px; font-weight: 700; letter-spacing: .5px; } +.version { font-size: 12px; color: var(--muted); font-weight: 400; margin-left: 8px; } +.header-right { display: flex; align-items: center; gap: 10px; } +.status-pill { + background: var(--panel2); border: 1px solid var(--border); + padding: 5px 12px; border-radius: 20px; font-size: 12px; color: var(--muted); +} +.status-pill.active { color: var(--green); border-color: var(--green); } +main { padding: 20px 24px; max-width: 1500px; margin: 0 auto; } + +/* ---------- buttons ---------- */ +.btn { + background: var(--panel2); color: var(--text); border: 1px solid var(--border); + padding: 7px 14px; border-radius: 8px; cursor: pointer; font-size: 13px; + transition: .15s; white-space: nowrap; +} +.btn:hover { border-color: var(--accent); color: #fff; } +.btn.primary { background: var(--accent); border-color: var(--accent); color: #fff; } +.btn.primary:hover { filter: brightness(1.1); } +.btn.ghost { background: transparent; } +.btn.danger:hover { border-color: var(--red); color: var(--red); } +.btn.sm { padding: 3px 8px; font-size: 12px; border-radius: 6px; } +.btn:disabled { opacity: .4; cursor: not-allowed; } + +/* ---------- task grid ---------- */ +.task-grid { display: grid; grid-template-columns: repeat(auto-fill, minmax(380px, 1fr)); gap: 16px; } +.card { + background: var(--panel); border: 1px solid var(--border); + border-radius: 12px; padding: 16px; display: flex; flex-direction: column; gap: 10px; + transition: .15s; +} +.card:hover { border-color: #3a4a75; transform: translateY(-1px); } +.card-head { display: flex; align-items: center; justify-content: space-between; gap: 8px; } +.card-name { font-size: 15px; font-weight: 600; overflow: hidden; text-overflow: ellipsis; white-space: nowrap; } +.card-meta { display: flex; flex-wrap: wrap; gap: 6px; align-items: center; } +.badge { + font-size: 11px; padding: 2px 8px; border-radius: 10px; + background: var(--panel2); border: 1px solid var(--border); color: var(--muted); +} +.badge.batch { color: #7cc4ff; border-color: #7cc4ff55; } +.badge.scheduled { color: var(--purple); border-color: var(--purple); } +.badge.auto { color: var(--yellow); border-color: var(--yellow); } +.badge.running { color: var(--green); border-color: var(--green); animation: pulse 1.6s infinite; } +.badge.paused { color: var(--yellow); border-color: var(--yellow); } +.badge.completed { color: var(--green); border-color: var(--green); } +.badge.failed { color: var(--red); border-color: var(--red); } +.badge.stopped { color: var(--muted); border-color: var(--border); } +.badge.scheduled-wait { color: var(--purple); border-color: var(--purple); } +@keyframes pulse { 50% { opacity: .5; } } + +.card-line { color: var(--muted); font-size: 12px; overflow: hidden; text-overflow: ellipsis; white-space: nowrap; } +.card-line b { color: var(--text); font-weight: 500; } +.progress { height: 6px; background: var(--panel2); border-radius: 4px; overflow: hidden; } +.progress > div { height: 100%; background: linear-gradient(90deg, var(--accent), var(--green)); transition: width .5s; } +.progress-label { display: flex; justify-content: space-between; font-size: 11px; color: var(--muted); } +.stats { display: flex; gap: 14px; font-size: 12px; color: var(--muted); } +.stats .ok { color: var(--green); } .stats .fail { color: var(--red); } .stats .img { color: var(--yellow); } +.card-actions { display: flex; gap: 6px; flex-wrap: wrap; margin-top: auto; } + +/* ---------- empty ---------- */ +.empty { text-align: center; padding: 80px 0; color: var(--muted); } +.empty-icon { font-size: 48px; margin-bottom: 12px; } +.empty button { margin-top: 14px; } + +/* ---------- modal ---------- */ +.modal-overlay { + position: fixed; inset: 0; background: rgba(5, 8, 15, .7); + display: flex; align-items: center; justify-content: center; z-index: 100; + backdrop-filter: blur(3px); +} +.modal { + background: var(--panel); border: 1px solid var(--border); border-radius: 14px; + width: 640px; max-width: 95vw; max-height: 92vh; display: flex; flex-direction: column; + box-shadow: 0 20px 60px rgba(0,0,0,.5); +} +.modal.wide { width: 1080px; } +.modal.tall { width: 1000px; height: 88vh; } +.modal-head { + display: flex; justify-content: space-between; align-items: center; + padding: 14px 18px; border-bottom: 1px solid var(--border); font-size: 15px; font-weight: 600; +} +.modal-foot { + display: flex; justify-content: flex-end; gap: 10px; align-items: center; + padding: 12px 18px; border-top: 1px solid var(--border); +} +.hint { color: var(--yellow); font-size: 12px; margin-right: auto; } + +.tabs { display: flex; gap: 6px; padding: 12px 18px 0; } +.tab { + padding: 8px 16px; background: transparent; border: 1px solid var(--border); + border-bottom: none; border-radius: 10px 10px 0 0; color: var(--muted); + cursor: pointer; font-size: 13px; +} +.tab.active { background: var(--panel2); color: #fff; border-color: var(--accent); } +.tab:disabled { opacity: .4; cursor: not-allowed; } + +.form-body { padding: 16px 18px; overflow-y: auto; display: flex; flex-direction: column; gap: 12px; } +.field { display: flex; flex-direction: column; gap: 5px; flex: 1; min-width: 0; } +.field label { font-size: 12px; color: var(--muted); } +.field.check { flex-direction: row; align-items: center; gap: 8px; padding-top: 22px; } +.field.check label { color: var(--text); font-size: 13px; display: flex; align-items: center; gap: 7px; cursor: pointer; } +input, select, textarea { + background: var(--panel2); border: 1px solid var(--border); color: var(--text); + border-radius: 8px; padding: 8px 10px; font-size: 13px; font-family: inherit; width: 100%; +} +input:focus, select:focus, textarea:focus { outline: none; border-color: var(--accent); } +textarea { resize: vertical; } +input[type="checkbox"] { width: auto; accent-color: var(--accent); } +.row2 { display: flex; gap: 12px; } +.box { border: 1px dashed var(--border); border-radius: 10px; padding: 12px; display: flex; flex-direction: column; gap: 10px; background: #141b2b; } +.hidden { display: none !important; } + +/* ---------- detail ---------- */ +.detail-wrap { padding: 16px 18px; overflow-y: auto; display: flex; flex-direction: column; gap: 14px; } +.detail-meta { display: flex; flex-wrap: wrap; gap: 14px; font-size: 12px; color: var(--muted); } +.detail-meta b { color: var(--text); font-weight: 500; } +.detail-actions { display: flex; gap: 8px; flex-wrap: wrap; } +.runs-bar { display: flex; gap: 8px; flex-wrap: wrap; align-items: center; } +.run-chip { + padding: 5px 12px; border-radius: 8px; border: 1px solid var(--border); + background: var(--panel2); cursor: pointer; font-size: 12px; color: var(--muted); +} +.run-chip.active { border-color: var(--accent); color: #fff; } +.run-panel { border: 1px solid var(--border); border-radius: 10px; padding: 14px; display: flex; flex-direction: column; gap: 10px; background: var(--panel2); } +.run-stats { display: flex; gap: 16px; font-size: 12px; flex-wrap: wrap; } +.table-wrap { overflow: auto; max-height: 320px; border: 1px solid var(--border); border-radius: 8px; } +table { width: 100%; border-collapse: collapse; font-size: 12px; } +th, td { padding: 7px 10px; text-align: left; border-bottom: 1px solid var(--border); white-space: nowrap; } +th { position: sticky; top: 0; background: var(--panel2); color: var(--muted); font-weight: 500; z-index: 1; } +td.url-cell { max-width: 300px; overflow: hidden; text-overflow: ellipsis; } +td.title-cell { max-width: 220px; overflow: hidden; text-overflow: ellipsis; } +.t-ok { color: var(--green); } .t-fail { color: var(--red); } +.file-link { color: var(--accent); cursor: pointer; text-decoration: underline; } +.file-link:hover { color: #9ec0ff; } +.thumbs { display: flex; flex-wrap: wrap; gap: 8px; } +.thumbs img { width: 96px; height: 72px; object-fit: cover; border-radius: 6px; border: 1px solid var(--border); cursor: pointer; } +.logs { + background: #0b0f1a; border: 1px solid var(--border); border-radius: 8px; + padding: 10px; font-family: ui-monospace, Consolas, monospace; font-size: 12px; + max-height: 220px; overflow-y: auto; line-height: 1.7; +} +.logs .info { color: var(--text); } .logs .warn { color: var(--yellow); } .logs .error { color: var(--red); } + +/* ---------- preview ---------- */ +#previewFrame { flex: 1; border: none; background: #fff; border-radius: 0 0 14px 14px; } + +/* ---------- misc ---------- */ +.toast { + position: fixed; top: 70px; left: 50%; transform: translateX(-50%); + background: var(--panel2); border: 1px solid var(--accent); color: #fff; + padding: 10px 20px; border-radius: 10px; z-index: 300; font-size: 13px; + box-shadow: 0 8px 30px rgba(0,0,0,.4); animation: fadein .2s; +} +.toast.error { border-color: var(--red); } +@keyframes fadein { from { opacity: 0; transform: translate(-50%, -8px); } } +::-webkit-scrollbar { width: 8px; height: 8px; } +::-webkit-scrollbar-thumb { background: #2a3550; border-radius: 4px; } diff --git a/store.py b/store.py new file mode 100644 index 0000000..0ff5060 --- /dev/null +++ b/store.py @@ -0,0 +1,122 @@ +# -*- coding: utf-8 -*- +"""任务 / 运行记录持久化层 (JSON 文件存储, 线程安全)""" +import json +import os +import threading +import uuid +from datetime import datetime + +HERE = os.path.dirname(os.path.abspath(__file__)) +DATA_DIR = os.path.join(HERE, "data") +TASKS_FILE = os.path.join(DATA_DIR, "tasks.json") +RUNS_FILE = os.path.join(DATA_DIR, "runs.json") # {task_id: [run, ...]} + +_lock = threading.RLock() + + +def new_id(prefix): + return f"{prefix}_{uuid.uuid4().hex[:12]}" + + +def now_str(): + return datetime.now().strftime("%Y-%m-%d %H:%M:%S") + + +def _load(path, default): + if not os.path.exists(path): + return default + try: + with open(path, encoding="utf-8") as f: + return json.load(f) + except Exception: + return default + + +def _save(path, data): + os.makedirs(os.path.dirname(path), exist_ok=True) + tmp = path + ".tmp" + with open(tmp, "w", encoding="utf-8") as f: + json.dump(data, f, ensure_ascii=False, indent=2) + os.replace(tmp, path) + + +# ---------------- 任务 ---------------- + +def load_tasks(): + with _lock: + return _load(TASKS_FILE, []) + + +def get_task(task_id): + for t in load_tasks(): + if t["id"] == task_id: + return t + return None + + +def upsert_task(task): + with _lock: + tasks = _load(TASKS_FILE, []) + for i, t in enumerate(tasks): + if t["id"] == task["id"]: + tasks[i] = task + break + else: + tasks.append(task) + _save(TASKS_FILE, tasks) + + +def delete_task(task_id): + with _lock: + tasks = _load(TASKS_FILE, []) + tasks = [t for t in tasks if t["id"] != task_id] + _save(TASKS_FILE, tasks) + runs = _load(RUNS_FILE, {}) + runs.pop(task_id, None) + _save(RUNS_FILE, runs) + + +# ---------------- 运行记录 ---------------- + +def load_runs_map(): + with _lock: + return _load(RUNS_FILE, {}) + + +def get_runs(task_id): + with _lock: + return _load(RUNS_FILE, {}).get(task_id, []) + + +def add_run(task_id, run): + with _lock: + runs = _load(RUNS_FILE, {}) + runs.setdefault(task_id, []).append(run) + if len(runs[task_id]) > 50: # 每个任务最多保留 50 次运行记录 + del runs[task_id][:-50] + _save(RUNS_FILE, runs) + + +def save_run(task_id, run): + with _lock: + runs = _load(RUNS_FILE, {}) + lst = runs.setdefault(task_id, []) + for i, r in enumerate(lst): + if r["id"] == run["id"]: + lst[i] = run + break + else: + lst.append(run) + if len(lst) > 50: + del lst[:-50] + _save(RUNS_FILE, runs) + + +def get_run(run_id): + with _lock: + runs = _load(RUNS_FILE, {}) + for lst in runs.values(): + for r in lst: + if r["id"] == run_id: + return r + return None