Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
143a3d5c89 | ||
|
|
57d45d901b | ||
|
|
950f5af5ea | ||
|
|
ec053266a7 | ||
|
|
32eebf3dd1 | ||
|
|
d9c9f0c633 | ||
|
|
909e8e01b5 | ||
|
|
e749d70a43 | ||
|
|
56f30f78c4 | ||
|
|
0482284cc0 | ||
|
|
796e667533 |
@@ -4,3 +4,5 @@ data/*.json
|
||||
data/cookies_*.json
|
||||
logs/
|
||||
out/
|
||||
data/exports/
|
||||
data/auto_state/
|
||||
@@ -4,14 +4,18 @@
|
||||
启动: /home/hz1/miniconda3/envs/openclaw/bin/python app.py (默认端口 16062)
|
||||
"""
|
||||
import os
|
||||
import re
|
||||
import threading
|
||||
import time
|
||||
import zipfile
|
||||
from datetime import datetime
|
||||
|
||||
from flask import Flask, jsonify, request, send_file, send_from_directory
|
||||
|
||||
import store
|
||||
import db
|
||||
from engine import CrawlJob, probe_links
|
||||
import notify
|
||||
from engine import CrawlJob, probe_links, normalize_url, url_excluded
|
||||
from scheduler import Scheduler, cron_next, interval_delta
|
||||
|
||||
HERE = os.path.dirname(os.path.abspath(__file__))
|
||||
@@ -95,23 +99,26 @@ def persist_cb(task_id, run):
|
||||
run["progress"]["percent"] = round(done * 100 / total) if total else 0
|
||||
store.save_run(task_id, run)
|
||||
try:
|
||||
db.sync_run(run, lambda tid, r: store.save_run(tid, r))
|
||||
db.sync_run_async(run) # 后台异步同步, 不再阻塞爬虫线程
|
||||
except Exception as e:
|
||||
print(f"[db] 同步失败: {e}", flush=True)
|
||||
|
||||
|
||||
def start_run(task):
|
||||
"""为任务启动一次爬取, 返回 (run, error)"""
|
||||
def start_run(task, skip_seed=False):
|
||||
"""为任务启动一次爬取, 返回 (run, error)
|
||||
skip_seed: auto 任务继续爬取模式, 跳过起始网址直接从缓存队列消费
|
||||
"""
|
||||
with JOBS_LOCK:
|
||||
job = JOBS.get(task["id"])
|
||||
if job and job.is_running():
|
||||
return None, "该任务已有正在运行的爬取"
|
||||
run = make_run(task)
|
||||
run["skip_seed"] = bool(skip_seed)
|
||||
store.add_run(task["id"], run)
|
||||
job = CrawlJob(task, run, persist_cb)
|
||||
JOBS[task["id"]] = job
|
||||
job.start()
|
||||
db.upsert_task(task) # 确保任务在库中
|
||||
db.upsert_task_async(task) # 确保任务在库中 (异步)
|
||||
return run, None
|
||||
|
||||
|
||||
@@ -147,12 +154,37 @@ def api_status():
|
||||
|
||||
# ---------------- API: 任务 ----------------
|
||||
|
||||
def _strip_auto_state(task):
|
||||
"""剥离任务对象中的待爬/已爬队列 (已独立存储, 避免大 JSON 拖慢接口)"""
|
||||
if task.get("mode") == "auto" and task.get("auto"):
|
||||
task["auto"] = {k: v for k, v in task["auto"].items()
|
||||
if k not in ("pending", "visited")}
|
||||
return task
|
||||
|
||||
|
||||
def _run_summary(run):
|
||||
"""运行记录摘要 (不含 results/logs), 供列表/详情页使用"""
|
||||
return {
|
||||
"id": run["id"],
|
||||
"task_id": run.get("task_id", ""),
|
||||
"mode": run.get("mode", ""),
|
||||
"status": run.get("status", ""),
|
||||
"progress": run.get("progress", {}),
|
||||
"started_at": run.get("started_at", ""),
|
||||
"finished_at": run.get("finished_at", ""),
|
||||
"stats": run.get("stats", {}),
|
||||
"out_dir": run.get("out_dir", ""),
|
||||
}
|
||||
|
||||
|
||||
@app.route("/api/tasks", methods=["GET"])
|
||||
def api_tasks():
|
||||
tasks = [t for t in store.load_tasks() if not t.get("deleted_at")]
|
||||
runs_map = store.load_runs_map() # 只全量读一次
|
||||
for t in tasks:
|
||||
runs = store.get_runs(t["id"])
|
||||
t["latest_run"] = runs[-1] if runs else None
|
||||
_strip_auto_state(t)
|
||||
runs = runs_map.get(t["id"], [])
|
||||
t["latest_run"] = _run_summary(runs[-1]) if runs else None
|
||||
with JOBS_LOCK:
|
||||
job = JOBS.get(t["id"])
|
||||
t["running"] = bool(job and job.is_running())
|
||||
@@ -220,7 +252,7 @@ def api_create_task():
|
||||
return jsonify({"error": "请至少填写一个网址"}), 400
|
||||
|
||||
store.upsert_task(task)
|
||||
db.upsert_task(task)
|
||||
db.upsert_task_async(task)
|
||||
return jsonify(task), 201
|
||||
|
||||
|
||||
@@ -229,11 +261,67 @@ 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)))
|
||||
# 分页参数: 指定 run + 页码 (默认最新 run 第 1 页)
|
||||
rid = request.args.get("run", "")
|
||||
try:
|
||||
page = max(int(request.args.get("page", 1)), 1)
|
||||
except ValueError:
|
||||
page = 1
|
||||
try:
|
||||
page_size = min(max(int(request.args.get("page_size", 100)), 10), 500)
|
||||
except ValueError:
|
||||
page_size = 100
|
||||
|
||||
runs = store.get_runs(tid)
|
||||
if rid:
|
||||
cur = next((r for r in runs if r["id"] == rid), None)
|
||||
if cur is None:
|
||||
cur = runs[-1] if runs else None # 指定的 run 不存在时回退最新
|
||||
else:
|
||||
cur = runs[-1] if runs else None
|
||||
with JOBS_LOCK:
|
||||
job = JOBS.get(tid)
|
||||
task["running"] = bool(job and job.is_running())
|
||||
return jsonify(task)
|
||||
|
||||
result = _strip_auto_state(dict(task))
|
||||
result["runs"] = [_run_summary(r) for r in reversed(runs)]
|
||||
result["run"] = None
|
||||
result["run_id"] = ""
|
||||
result["run_total"] = 0
|
||||
result["run_page"] = 1
|
||||
result["run_pages"] = 1
|
||||
if cur:
|
||||
results = cur.get("results", [])
|
||||
total = len(results)
|
||||
pages = max((total + page_size - 1) // page_size, 1)
|
||||
page = min(page, pages)
|
||||
start = (page - 1) * page_size
|
||||
snap = {k: v for k, v in cur.items() if k != "logs"} # logs 走独立轮询接口
|
||||
snap["results"] = results[start:start + page_size]
|
||||
snap["results_total"] = total
|
||||
snap["run_page"] = page
|
||||
snap["run_pages"] = pages
|
||||
snap["run_page_size"] = page_size
|
||||
result["run"] = snap
|
||||
result["run_id"] = cur["id"]
|
||||
result["run_total"] = total
|
||||
result["run_page"] = page
|
||||
result["run_pages"] = pages
|
||||
result["run_page_size"] = page_size
|
||||
# auto 状态数量: 运行中优先读内存实时值, 否则读状态文件
|
||||
if task.get("mode") == "auto":
|
||||
live_v = live_p = None
|
||||
with JOBS_LOCK:
|
||||
job = JOBS.get(tid)
|
||||
if job:
|
||||
live_v = getattr(job, "_auto_visited", None)
|
||||
live_p = getattr(job, "_auto_pending", None)
|
||||
st = store.load_auto_state(tid, task)
|
||||
result["auto_pending_count"] = (
|
||||
live_p if live_p is not None else len(st.get("pending", [])))
|
||||
result["auto_visited_count"] = (
|
||||
live_v if live_v is not None else len(st.get("visited", [])))
|
||||
return jsonify(result)
|
||||
|
||||
|
||||
@app.route("/api/tasks/<tid>", methods=["PUT"])
|
||||
@@ -257,7 +345,22 @@ def api_update_task(tid):
|
||||
if running:
|
||||
job.update_config(body["config"]) # 运行中热更新
|
||||
if "auto" in body and task.get("mode") == "auto":
|
||||
old_exclude = task.get("auto", {}).get("exclude") or []
|
||||
task["auto"] = {**task.get("auto", {}), **body["auto"]}
|
||||
new_exclude = task["auto"].get("exclude") or []
|
||||
# 排除规则发生变化时, 同步清理待爬队列中已命中的 URL (已爬 visited 保留)
|
||||
if new_exclude and new_exclude != old_exclude and not running:
|
||||
st = store.load_auto_state(tid, task)
|
||||
kept = [p for p in st.get("pending", [])
|
||||
if not url_excluded(p.get("url", ""), task["auto"])]
|
||||
removed = len(st.get("pending", [])) - len(kept)
|
||||
if removed:
|
||||
store.save_auto_state(tid, {
|
||||
"pending": kept,
|
||||
"visited": st.get("visited", []),
|
||||
})
|
||||
if "exclude" in body.get("auto", {}):
|
||||
task["auto"]["_cleaned"] = removed
|
||||
if "schedule" in body and task.get("mode") == "scheduled":
|
||||
sch = {**task.get("schedule", {}), **body["schedule"]}
|
||||
try:
|
||||
@@ -267,7 +370,7 @@ def api_update_task(tid):
|
||||
task["schedule"] = sch
|
||||
task["updated_at"] = now_str()
|
||||
store.upsert_task(task)
|
||||
db.upsert_task(task)
|
||||
db.upsert_task_async(task)
|
||||
return jsonify(task)
|
||||
|
||||
|
||||
@@ -283,7 +386,7 @@ def api_delete_task(tid):
|
||||
with JOBS_LOCK:
|
||||
JOBS.pop(tid, None)
|
||||
store.soft_delete_task(tid, now_str())
|
||||
db.upsert_task(store.get_task(tid))
|
||||
db.upsert_task_async(store.get_task(tid))
|
||||
return jsonify({"ok": True, "msg": "已移入回收站"})
|
||||
|
||||
|
||||
@@ -306,8 +409,9 @@ def _purge_out_dir(out_dir):
|
||||
@app.route("/api/trash", methods=["GET"])
|
||||
def api_trash_list():
|
||||
items = store.list_trash()
|
||||
runs_map = store.load_runs_map()
|
||||
for t in items:
|
||||
t["runs_count"] = len(store.get_runs(t["id"]))
|
||||
t["runs_count"] = len(runs_map.get(t["id"], []))
|
||||
t["out_dir"] = resolve_out_dir(t)
|
||||
return jsonify(items)
|
||||
|
||||
@@ -318,7 +422,7 @@ def api_trash_restore(tid):
|
||||
if not task or not task.get("deleted_at"):
|
||||
return jsonify({"error": "任务不在回收站中"}), 404
|
||||
store.restore_task(tid)
|
||||
db.upsert_task(store.get_task(tid))
|
||||
db.upsert_task_async(store.get_task(tid))
|
||||
return jsonify({"ok": True, "msg": "已恢复"})
|
||||
|
||||
|
||||
@@ -346,6 +450,71 @@ def api_trash_clear():
|
||||
|
||||
# ---------------- API: 运行控制 ----------------
|
||||
|
||||
@app.route("/api/tasks/<tid>/retry-failed", methods=["POST"])
|
||||
def api_retry_failed(tid):
|
||||
"""重爬失败页: 提取指定 run (默认最新) 中失败的 URL, 从已爬集合解除标记并注入
|
||||
待爬队列头部; 返回注入数量, 调用方可随后点「继续爬取」重爬这些页面"""
|
||||
task = store.get_task(tid)
|
||||
if not task:
|
||||
return jsonify({"error": "任务不存在"}), 404
|
||||
if task.get("mode") != "auto":
|
||||
return jsonify({"error": "仅自动模式任务支持重爬失败页"}), 400
|
||||
try:
|
||||
body = request.get_json(force=True) or {}
|
||||
except Exception:
|
||||
body = {}
|
||||
rid = request.args.get("run", "") or body.get("run", "")
|
||||
runs = store.get_runs(tid)
|
||||
cur = None
|
||||
if rid:
|
||||
cur = next((r for r in runs if r["id"] == rid), None)
|
||||
else:
|
||||
cur = runs[-1] if runs else None
|
||||
if not cur:
|
||||
return jsonify({"error": "运行记录不存在"}), 404
|
||||
failed = [res.get("url") for res in cur.get("results", [])
|
||||
if res.get("status") == "FAIL" and res.get("url")]
|
||||
if not failed:
|
||||
return jsonify({"error": "该运行记录没有失败页面", "injected": 0})
|
||||
st = store.load_auto_state(tid, task)
|
||||
visited = set(st.get("visited", []))
|
||||
pending = st.get("pending", [])
|
||||
pending_urls = {normalize_url(p.get("url", "")) for p in pending}
|
||||
# 待重爬: 已爬过且不在待爬队列中的失败 URL (去重)
|
||||
to_inject, seen = [], set()
|
||||
for u in failed:
|
||||
key = normalize_url(u)
|
||||
if key in visited and key not in pending_urls and key not in seen:
|
||||
seen.add(key)
|
||||
to_inject.append({"url": u, "depth": 0, "source": "retry-failed"})
|
||||
if not to_inject:
|
||||
return jsonify({"error": "失败页面均已爬或已在待爬队列中", "injected": 0})
|
||||
# 解除已爬标记
|
||||
remove_keys = {normalize_url(u) for u in failed}
|
||||
visited = {v for v in visited if normalize_url(v) not in remove_keys}
|
||||
# 注入队列头部, 优先重爬
|
||||
pending = to_inject + pending
|
||||
store.save_auto_state(tid, {"pending": pending, "visited": sorted(visited)})
|
||||
return jsonify({"injected": len(to_inject), "pending_total": len(pending)})
|
||||
|
||||
|
||||
@app.route("/api/tasks/<tid>/continue", methods=["POST"])
|
||||
def api_continue(tid):
|
||||
"""继续爬取: auto 任务从待爬缓存队列接着爬 (跳过起始网址, 保留已爬集合)"""
|
||||
task = store.get_task(tid)
|
||||
if not task:
|
||||
return jsonify({"error": "任务不存在"}), 404
|
||||
if task.get("mode") != "auto":
|
||||
return jsonify({"error": "仅自动爬取任务支持继续爬取"}), 400
|
||||
st = store.load_auto_state(tid, task)
|
||||
if not st.get("pending"):
|
||||
return jsonify({"error": "没有待爬缓存链接,无需继续"}), 400
|
||||
run, err = start_run(task, skip_seed=True)
|
||||
if err:
|
||||
return jsonify({"error": err}), 409
|
||||
return jsonify(run)
|
||||
|
||||
|
||||
@app.route("/api/tasks/<tid>/start", methods=["POST"])
|
||||
def api_start(tid):
|
||||
task = store.get_task(tid)
|
||||
@@ -386,13 +555,20 @@ def api_resume(tid):
|
||||
|
||||
# ---------------- API: 统计 ----------------
|
||||
|
||||
# 磁盘占用统计缓存 (避免每次轮询都逐个 stat 文件)
|
||||
_DISK_CACHE = {"ts": 0.0, "mb": 0.0}
|
||||
DISK_CACHE_TTL = 30 # 秒
|
||||
|
||||
|
||||
@app.route("/api/stats")
|
||||
def api_stats():
|
||||
tasks = [t for t in store.load_tasks() if not t.get("deleted_at")]
|
||||
trash_count = sum(1 for t in store.load_tasks() if t.get("deleted_at"))
|
||||
tasks = store.load_tasks()
|
||||
active = [t for t in tasks if not t.get("deleted_at")]
|
||||
trash_count = len(tasks) - len(active)
|
||||
runs_map = store.load_runs_map() # 只全量读一次
|
||||
total_runs = ok = fail = imgs = 0
|
||||
for t in tasks:
|
||||
for r in store.get_runs(t["id"]):
|
||||
for t in active:
|
||||
for r in runs_map.get(t["id"], []):
|
||||
total_runs += 1
|
||||
st = r.get("stats") or {}
|
||||
ok += st.get("ok", 0)
|
||||
@@ -400,28 +576,30 @@ def api_stats():
|
||||
imgs += st.get("images", 0)
|
||||
with JOBS_LOCK:
|
||||
running = sum(1 for j in JOBS.values() if j.is_running())
|
||||
# 统计各任务输出目录的磁盘占用
|
||||
size = 0
|
||||
seen = set()
|
||||
for t in tasks:
|
||||
d = os.path.realpath(resolve_out_dir(t))
|
||||
if d in seen or not os.path.isdir(d):
|
||||
continue
|
||||
seen.add(d)
|
||||
for root, _dirs, files in os.walk(d):
|
||||
for f in files:
|
||||
try:
|
||||
size += os.path.getsize(os.path.join(root, f))
|
||||
except OSError:
|
||||
pass
|
||||
now = time.time()
|
||||
if now - _DISK_CACHE["ts"] > DISK_CACHE_TTL:
|
||||
size = 0
|
||||
seen = set()
|
||||
for t in active:
|
||||
d = os.path.realpath(resolve_out_dir(t))
|
||||
if d in seen or not os.path.isdir(d):
|
||||
continue
|
||||
seen.add(d)
|
||||
for root, _dirs, files in os.walk(d):
|
||||
for f in files:
|
||||
try:
|
||||
size += os.path.getsize(os.path.join(root, f))
|
||||
except OSError:
|
||||
pass
|
||||
_DISK_CACHE.update(ts=now, mb=round(size / 1048576, 1))
|
||||
return jsonify({
|
||||
"tasks": len(tasks),
|
||||
"tasks": len(active),
|
||||
"running": running,
|
||||
"runs": total_runs,
|
||||
"ok": ok,
|
||||
"fail": fail,
|
||||
"images": imgs,
|
||||
"disk_mb": round(size / 1048576, 1),
|
||||
"disk_mb": _DISK_CACHE["mb"],
|
||||
"trash": trash_count,
|
||||
})
|
||||
|
||||
@@ -463,9 +641,10 @@ def api_clear_cache(tid):
|
||||
auto = task.setdefault("auto", {})
|
||||
auto["pending"] = []
|
||||
auto["visited"] = []
|
||||
store.clear_auto_state(tid) # 同时清空独立状态文件
|
||||
task["updated_at"] = now_str()
|
||||
store.upsert_task(task)
|
||||
db.upsert_task(task)
|
||||
db.upsert_task_async(task)
|
||||
return jsonify({"ok": True, "msg": "缓存队列已清空"})
|
||||
|
||||
|
||||
@@ -500,6 +679,7 @@ def api_search():
|
||||
ql = q.lower()
|
||||
results = []
|
||||
seen = set()
|
||||
runs_map = store.load_runs_map() # 只全量读一次
|
||||
for task in store.load_tasks():
|
||||
if task.get("deleted_at"):
|
||||
continue # 回收站任务不参与搜索
|
||||
@@ -509,7 +689,7 @@ def api_search():
|
||||
or any(ql in u.lower() for u in task.get("urls", []))
|
||||
or ql in ((auto.get("seed_url") or "").lower())
|
||||
)
|
||||
for run in store.get_runs(task["id"]):
|
||||
for run in runs_map.get(task["id"], []):
|
||||
for r in run.get("results", []):
|
||||
url = r.get("url") or ""
|
||||
title = r.get("title") or ""
|
||||
@@ -588,18 +768,137 @@ def api_file():
|
||||
return send_file(full)
|
||||
|
||||
|
||||
# ---------------- API: 打包导出 ----------------
|
||||
|
||||
EXPORT_DIR = os.path.join(HERE, "data", "exports")
|
||||
EXPORT_KEEP = 10 # 最多保留的打包文件数
|
||||
|
||||
|
||||
def _safe_zip_name(task):
|
||||
"""任务名安全化为文件名 (保留中文, 去非法字符)"""
|
||||
name = re.sub(r'[\\/:*?"<>|\s]+', "_", task.get("name", "")).strip("_")
|
||||
return (name or task["id"])[:60]
|
||||
|
||||
|
||||
def _make_zip(task):
|
||||
"""把任务输出目录打包为 zip, 返回 (zip_path, err) 或 (None, 错误信息)"""
|
||||
out_dir = os.path.realpath(resolve_out_dir(task))
|
||||
if not os.path.isdir(out_dir):
|
||||
return None, "输出目录不存在"
|
||||
os.makedirs(EXPORT_DIR, exist_ok=True)
|
||||
zip_path = os.path.join(EXPORT_DIR, f"{task['id']}_{int(time.time())}.zip")
|
||||
prefix = _safe_zip_name(task) + "/"
|
||||
count = 0
|
||||
try:
|
||||
with zipfile.ZipFile(zip_path, "w", zipfile.ZIP_DEFLATED) as zf:
|
||||
for root, _dirs, files in os.walk(out_dir):
|
||||
for f in files:
|
||||
full = os.path.join(root, f)
|
||||
rel = os.path.relpath(full, out_dir)
|
||||
zf.write(full, prefix + rel)
|
||||
count += 1
|
||||
except Exception as e:
|
||||
try:
|
||||
os.remove(zip_path)
|
||||
except OSError:
|
||||
pass
|
||||
return None, f"打包失败: {e}"
|
||||
_cleanup_exports()
|
||||
return zip_path, None
|
||||
|
||||
|
||||
def _cleanup_exports():
|
||||
"""清理旧的打包文件, 只保留最近 EXPORT_KEEP 个"""
|
||||
try:
|
||||
files = sorted(
|
||||
(os.path.join(EXPORT_DIR, f) for f in os.listdir(EXPORT_DIR)
|
||||
if f.endswith(".zip")),
|
||||
key=os.path.getmtime, reverse=True,
|
||||
)
|
||||
for f in files[EXPORT_KEEP:]:
|
||||
os.remove(f)
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
|
||||
@app.route("/api/export")
|
||||
def api_export():
|
||||
"""打包任务输出目录为 zip 并提供下载"""
|
||||
tid = request.args.get("task_id", "")
|
||||
task = store.get_task(tid)
|
||||
if not task:
|
||||
return jsonify({"error": "任务不存在"}), 404
|
||||
zip_path, err = _make_zip(task)
|
||||
if err:
|
||||
return jsonify({"error": err}), 400
|
||||
return send_file(
|
||||
zip_path, as_attachment=True,
|
||||
download_name=_safe_zip_name(task) + ".zip",
|
||||
mimetype="application/zip",
|
||||
)
|
||||
|
||||
|
||||
@app.route("/api/export/email", methods=["POST"])
|
||||
def api_export_email():
|
||||
"""打包任务输出目录为 zip 并发送到指定邮箱 (默认 wlq@tphai.com)"""
|
||||
body = request.get_json(force=True) or {}
|
||||
tid = body.get("task_id", "")
|
||||
email = (body.get("email") or "").strip() or "wlq@tphai.com"
|
||||
task = store.get_task(tid)
|
||||
if not task:
|
||||
return jsonify({"error": "任务不存在"}), 404
|
||||
if "@" not in email:
|
||||
return jsonify({"error": "邮箱格式不正确"}), 400
|
||||
zip_path, err = _make_zip(task)
|
||||
if err:
|
||||
return jsonify({"error": err}), 400
|
||||
size_mb = round(os.path.getsize(zip_path) / 1048576, 1)
|
||||
body_text = (f"项目: {task['name']}\n"
|
||||
f"输出目录: {resolve_out_dir(task)}\n"
|
||||
f"打包文件: {os.path.basename(zip_path)}\n"
|
||||
f"压缩包大小: {size_mb} MB\n\n"
|
||||
f"打包时间: {now_str()}")
|
||||
ok, msg = notify.send_attachment(
|
||||
f"[爬虫打包] {task['name']}", body_text, email, zip_path)
|
||||
if not ok:
|
||||
return jsonify({"error": f"邮件发送失败: {msg}"}), 500
|
||||
return jsonify({"ok": True, "msg": f"已发送到 {email} (zip {size_mb} MB)"})
|
||||
|
||||
|
||||
# ---------------- 启动 ----------------
|
||||
|
||||
scheduler = Scheduler(start_run)
|
||||
|
||||
|
||||
def migrate_auto_state():
|
||||
"""启动时把旧任务对象里的 auto pending/visited 迁移到独立文件, 给 tasks.json 瘦身"""
|
||||
migrated = 0
|
||||
for t in store.load_tasks():
|
||||
if t.get("mode") != "auto":
|
||||
continue
|
||||
auto = t.get("auto") or {}
|
||||
if auto.get("pending") or auto.get("visited"):
|
||||
store.save_auto_state(t["id"], {
|
||||
"pending": auto.get("pending", []),
|
||||
"visited": auto.get("visited", []),
|
||||
})
|
||||
auto.pop("pending", None)
|
||||
auto.pop("visited", None)
|
||||
t["updated_at"] = now_str()
|
||||
store.upsert_task(t)
|
||||
migrated += 1
|
||||
if migrated:
|
||||
print(f"[universal-crawler] 已迁移 {migrated} 个 auto 任务的队列状态到独立文件", flush=True)
|
||||
|
||||
|
||||
def main():
|
||||
os.makedirs(os.path.join(HERE, "data"), exist_ok=True)
|
||||
os.makedirs(os.path.join(HERE, "out"), exist_ok=True)
|
||||
db.init_db()
|
||||
# 全量同步存量任务到数据库
|
||||
migrate_auto_state()
|
||||
# 全量同步存量任务到数据库 (异步, 不阻塞启动)
|
||||
for t in store.load_tasks():
|
||||
db.upsert_task(t)
|
||||
db.upsert_task_async(t)
|
||||
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)
|
||||
|
||||
@@ -5,9 +5,11 @@ MySQL 记录层: 任务 / 运行 / 爬取结果 写入数据库
|
||||
- 爬取成功/失败均有 status 标记 (OK / FAIL), 运行记录有 run 状态标记
|
||||
- 所有操作容错: 数据库不可用时不影响爬取主流程 (仅打印日志)
|
||||
"""
|
||||
import copy
|
||||
import json
|
||||
import os
|
||||
import time
|
||||
from concurrent.futures import ThreadPoolExecutor
|
||||
|
||||
import pymysql
|
||||
|
||||
@@ -101,6 +103,39 @@ def _safe(fn, *args, **kwargs):
|
||||
print(f"[db] 操作失败: {e}", flush=True)
|
||||
|
||||
|
||||
# ---------------- 异步同步 (避免远程 MySQL 阻塞爬虫/API 主流程) ----------------
|
||||
|
||||
_db_executor = ThreadPoolExecutor(max_workers=2, thread_name_prefix="db-sync")
|
||||
|
||||
|
||||
def upsert_task_async(task):
|
||||
"""异步写入/更新任务 (不阻塞调用方)"""
|
||||
try:
|
||||
snap = copy.deepcopy(task)
|
||||
_db_executor.submit(upsert_task, snap)
|
||||
except Exception as e:
|
||||
print(f"[db] 异步任务同步失败: {e}", flush=True)
|
||||
|
||||
|
||||
def sync_run_async(run):
|
||||
"""异步同步运行记录+爬取结果; 完成后把 _db_count 回写到主线程 run 对象"""
|
||||
try:
|
||||
snap = copy.deepcopy(run)
|
||||
|
||||
def _work():
|
||||
upsert_run(snap)
|
||||
results = snap.get("results", [])
|
||||
synced = snap.get("_db_count", 0)
|
||||
if len(results) > synced:
|
||||
insert_results(snap, results[synced:])
|
||||
snap["_db_count"] = len(results)
|
||||
run["_db_count"] = snap["_db_count"]
|
||||
|
||||
_db_executor.submit(_work)
|
||||
except Exception as e:
|
||||
print(f"[db] 异步运行同步失败: {e}", flush=True)
|
||||
|
||||
|
||||
def init_db():
|
||||
"""建库建表 (幂等)"""
|
||||
try:
|
||||
|
||||
@@ -30,21 +30,60 @@ 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 = [
|
||||
# title 中出现任一标记即判定反爬 (title 短小可靠, 不会误伤正文)
|
||||
_CHALLENGE_TITLE_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",
|
||||
]
|
||||
|
||||
# html 中仅匹配强特征标记, 且只在 head 区域(前 20KB)搜索。
|
||||
# 不能全文匹配 "captcha"/"challenge" 等词——正常文章正文提到这些词会被误判为反爬。
|
||||
_CHALLENGE_HTML_MARKS = [
|
||||
"cf-challenge", "challenge-platform", "cf-browser-verification",
|
||||
"verify you are human", "checking your browser",
|
||||
"enable javascript and cookies", "just a moment",
|
||||
]
|
||||
|
||||
# 反爬拦截类错误信号: 命中后不再重试 (重试无意义且拖慢任务)
|
||||
_ANTI_CRAWL_MARKS = [
|
||||
"403", "429", "503", "forbidden", "access denied", "too many requests",
|
||||
"captcha", "cloudflare", "challenge", "verify you are human",
|
||||
"just a moment", "blocked", "被反爬拦截",
|
||||
]
|
||||
|
||||
|
||||
def is_anti_crawl_error(err):
|
||||
"""判断错误是否属于反爬拦截 (此类失败重试也无法通过, 直接放弃该页)"""
|
||||
s = str(err or "").lower()
|
||||
return any(m in s for m in _ANTI_CRAWL_MARKS)
|
||||
|
||||
|
||||
_BROWSER_DEAD_MARKS = [
|
||||
"browser has been closed", "target page, context or browser has been closed",
|
||||
"has been disposed", "execution context was destroyed", "browser closed",
|
||||
"page closed", "target closed", "crash",
|
||||
]
|
||||
|
||||
|
||||
def is_browser_dead_error(err):
|
||||
"""判断错误是否属于浏览器/页面失效 (需重启浏览器后重试, 重试同一 URL 才有意义)"""
|
||||
s = str(err or "").lower()
|
||||
return any(m in s for m in _BROWSER_DEAD_MARKS)
|
||||
|
||||
|
||||
def is_challenge_page(title, html):
|
||||
"""判断是否仍在反爬验证页 (仅关键词启发, 避免误伤正常小页面)"""
|
||||
low = html.lower()
|
||||
"""判断是否仍在反爬验证页
|
||||
- title 出现反爬关键词: 判定 (可靠, 反爬页 title 基本都会变)
|
||||
- html 仅在前 20KB(head 区域)出现强特征标记时判定,
|
||||
避免正文含 captcha/challenge 等词的正常页面被误杀
|
||||
"""
|
||||
t = (title or "").lower()
|
||||
for mark in _CHALLENGE_MARKS:
|
||||
if mark in t or mark in low:
|
||||
for mark in _CHALLENGE_TITLE_MARKS:
|
||||
if mark in t:
|
||||
return True
|
||||
return False
|
||||
low = (html or "").lower()[:20000]
|
||||
return any(m in low for m in _CHALLENGE_HTML_MARKS)
|
||||
|
||||
|
||||
def safe_name(url, idx):
|
||||
@@ -56,6 +95,7 @@ def safe_name(url, idx):
|
||||
def _settle_wait(page, timeout_s):
|
||||
"""等待页面稳定(无 stop/pause 检查, 供试爬取等独立流程使用)"""
|
||||
last_title, stable = "", 0
|
||||
challenge_hits = 0
|
||||
start = time.time()
|
||||
while time.time() - start < timeout_s:
|
||||
time.sleep(1)
|
||||
@@ -65,8 +105,12 @@ def _settle_wait(page, timeout_s):
|
||||
except Exception:
|
||||
continue
|
||||
if is_challenge_page(title, html):
|
||||
challenge_hits += 1
|
||||
if challenge_hits >= 8: # 持续 8 秒仍是验证页, 判定反爬, 尽早放弃
|
||||
return True, title, html
|
||||
stable = 0
|
||||
continue
|
||||
challenge_hits = 0
|
||||
if title == last_title:
|
||||
stable += 1
|
||||
if stable >= 2 and len(html) > 1000:
|
||||
@@ -107,6 +151,21 @@ def normalize_url(url):
|
||||
return str(url)
|
||||
|
||||
|
||||
def url_excluded(url, auto=None):
|
||||
"""判断 URL 是否命中排除规则 (与 filter_links 的 exclude 判定逻辑一致, 供队列清理使用)"""
|
||||
if not auto:
|
||||
return False
|
||||
exclude = auto.get("exclude") or []
|
||||
if not exclude:
|
||||
return False
|
||||
use_regex = bool(auto.get("use_regex"))
|
||||
s = str(url or "")
|
||||
if use_regex:
|
||||
return any(re.search(p, s) for p in exclude)
|
||||
low = s.lower()
|
||||
return any(p.lower() in low for p in exclude)
|
||||
|
||||
|
||||
def filter_links(hrefs, seed_url, include=None, exclude=None,
|
||||
same_domain=True, use_regex=False):
|
||||
"""按规则过滤链接, 返回 (included, excluded); excluded 含排除原因"""
|
||||
@@ -208,8 +267,19 @@ class CrawlJob:
|
||||
self.persist = persist # callable(task_id, run)
|
||||
self._stop = threading.Event()
|
||||
self._pause = threading.Event()
|
||||
# 运行中的 auto 实时状态 (供详情接口读取; 任务结束时由状态文件兜底)
|
||||
self._auto_visited = None
|
||||
self._auto_pending = None
|
||||
self._cfg_lock = threading.RLock()
|
||||
self.thread = None
|
||||
# 浏览器句柄 (p, browser, ctx, page, cookie_file); 崩溃后重建
|
||||
self._browser = None
|
||||
|
||||
def _page(self):
|
||||
return self._browser[3] if self._browser else None
|
||||
|
||||
def _ctx(self):
|
||||
return self._browser[2] if self._browser else None
|
||||
|
||||
# ---------------- 控制接口 ----------------
|
||||
def start(self):
|
||||
@@ -288,9 +358,13 @@ class CrawlJob:
|
||||
pass
|
||||
Stealth().apply_stealth_sync(ctx)
|
||||
page = ctx.new_page()
|
||||
return p, browser, ctx, page, cookie_file
|
||||
self._browser = (p, browser, ctx, page, cookie_file)
|
||||
return self._browser
|
||||
|
||||
def _close_browser(self, p, browser, ctx, cookie_file):
|
||||
def _close_browser(self):
|
||||
if not self._browser:
|
||||
return
|
||||
p, browser, ctx, _page, cookie_file = self._browser
|
||||
try:
|
||||
json.dump(ctx.cookies(), open(cookie_file, "w"))
|
||||
except Exception:
|
||||
@@ -303,9 +377,24 @@ class CrawlJob:
|
||||
p.stop()
|
||||
except Exception:
|
||||
pass
|
||||
self._browser = None
|
||||
|
||||
def _reopen_browser(self):
|
||||
"""浏览器失效后重建: 保存 cookie -> 关闭旧实例 -> 启动新实例"""
|
||||
self._close_browser()
|
||||
time.sleep(1)
|
||||
self._open_browser()
|
||||
self._log("info", "浏览器已重启")
|
||||
|
||||
def _need_browser_restart(self, done_count):
|
||||
"""每爬 N 页主动重启一次浏览器, 防止长时间运行内存膨胀导致崩溃 (browser_max_pages=0 关闭)"""
|
||||
limit = int(self._cfg("browser_max_pages", 0) or 0)
|
||||
return limit > 0 and done_count > 1 and (done_count - 1) % limit == 0
|
||||
|
||||
def _wait_page_settle(self, page, timeout_s):
|
||||
last_title, stable = "", 0
|
||||
challenge_hits = 0
|
||||
err_streak = 0 # 连续读取失败计数
|
||||
start = time.time()
|
||||
while time.time() - start < timeout_s:
|
||||
if self._stop.is_set():
|
||||
@@ -316,10 +405,18 @@ class CrawlJob:
|
||||
title = page.title()
|
||||
html = page.content()
|
||||
except Exception:
|
||||
err_streak += 1
|
||||
if err_streak >= 3: # 页面/浏览器已失效, 提前失败而非干等到超时
|
||||
raise RuntimeError("Target page, context or browser has been closed")
|
||||
continue # 正在跳转
|
||||
err_streak = 0
|
||||
if is_challenge_page(title, html):
|
||||
challenge_hits += 1
|
||||
if challenge_hits >= 8: # 持续 8 秒仍是验证页, 判定反爬, 尽早放弃
|
||||
return True, title, html
|
||||
stable = 0
|
||||
continue
|
||||
challenge_hits = 0
|
||||
if title == last_title:
|
||||
stable += 1
|
||||
if stable >= 2 and len(html) > 1000:
|
||||
@@ -435,7 +532,9 @@ class CrawlJob:
|
||||
"source_url": source_url, "depth": depth,
|
||||
"images": [], "attempts": 0,
|
||||
}
|
||||
for attempt in range(retries + 1):
|
||||
dead_strikes = 0 # 浏览器失效重建次数 (防止无限重建)
|
||||
attempt = 0
|
||||
while attempt <= retries:
|
||||
if self._stop.is_set():
|
||||
entry["error"] = "任务已终止"
|
||||
break
|
||||
@@ -459,7 +558,28 @@ class CrawlJob:
|
||||
except Exception as e:
|
||||
entry["error"] = str(e)
|
||||
entry["crawl_time"] = store.now_str()
|
||||
if is_browser_dead_error(e):
|
||||
# 浏览器/页面失效: 重启浏览器后重试同一 URL, 不消耗重试次数
|
||||
dead_strikes += 1
|
||||
if dead_strikes > 3:
|
||||
entry["error"] = f"浏览器多次重启仍失效, 放弃: {e}"
|
||||
self._log("error", f"浏览器多次重启仍失效, 放弃 {url}")
|
||||
break
|
||||
self._log("warn", f"浏览器已失效({e}), 重启后重试: {url}")
|
||||
try:
|
||||
self._reopen_browser()
|
||||
page, ctx = self._page(), self._ctx()
|
||||
except Exception as re_err:
|
||||
entry["error"] = f"浏览器重启失败: {re_err}"
|
||||
self._log("error", f"浏览器重启失败, 放弃 {url}: {re_err}")
|
||||
break
|
||||
continue
|
||||
self._log("warn", f"第{attempt + 1}次失败 {url}: {e}")
|
||||
if is_anti_crawl_error(e):
|
||||
# 反爬拦截: 重试也过不去, 直接放弃, 不再消耗重试次数
|
||||
entry["error"] = f"反爬拦截, 跳过重试: {e}"
|
||||
self._log("warn", f"判定为反爬拦截, 放弃重试: {url}")
|
||||
break
|
||||
if attempt < retries:
|
||||
self._wait_if_paused()
|
||||
t0 = time.time()
|
||||
@@ -468,6 +588,7 @@ class CrawlJob:
|
||||
break
|
||||
self._wait_if_paused()
|
||||
time.sleep(0.3)
|
||||
attempt += 1
|
||||
self._write_page_meta(out_dir, base, entry)
|
||||
return entry
|
||||
|
||||
@@ -531,17 +652,20 @@ class CrawlJob:
|
||||
os.makedirs(out_dir, exist_ok=True)
|
||||
self._persist()
|
||||
|
||||
p, browser, ctx, page, cookie_file = self._open_browser()
|
||||
self._open_browser()
|
||||
try:
|
||||
for i, url in enumerate(urls, 1):
|
||||
if self._stop.is_set():
|
||||
self._log("info", "收到终止信号, 停止爬取")
|
||||
break
|
||||
self._wait_if_paused()
|
||||
if self._need_browser_restart(i):
|
||||
self._log("info", f"已爬 {i - 1} 页, 主动重启浏览器")
|
||||
self._reopen_browser()
|
||||
run["progress"]["current_url"] = url
|
||||
run["progress"]["done"] = i - 1
|
||||
self._persist()
|
||||
entry = self._retry_crawl(page, ctx, url, i, out_dir)
|
||||
entry = self._retry_crawl(self._page(), self._ctx(), url, i, out_dir)
|
||||
run["results"].append(entry)
|
||||
self._bump_stats(entry)
|
||||
run["progress"]["done"] = i
|
||||
@@ -549,7 +673,7 @@ class CrawlJob:
|
||||
if entry["status"] == "OK":
|
||||
self._delay()
|
||||
finally:
|
||||
self._close_browser(p, browser, ctx, cookie_file)
|
||||
self._close_browser()
|
||||
|
||||
def _discover_links(self, page):
|
||||
"""从当前页面提取符合规则的链接"""
|
||||
@@ -570,9 +694,10 @@ class CrawlJob:
|
||||
def _crawl_auto(self):
|
||||
"""自动爬取: 持续递归 爬取->发现链接->爬取... 直到无新链接可爬或达到上限
|
||||
- max_depth: 0=无限制, N=只爬 N 层
|
||||
- max_pages: 0=无限制, N=安全上限
|
||||
- max_pages: 0=无限制, N=安全上限 (继续爬取模式按本次新增页数重新计算)
|
||||
- 提取但未爬取的链接持久化到 auto.pending, 已爬集合持久化到 auto.visited
|
||||
- 停止后再次运行从缓存队列继续爬; 起始网址每次运行都重新爬(不去重)
|
||||
- skip_seed(继续爬取): 跳过起始网址直接消费缓存队列, 页数上限按本次新增重新计算
|
||||
"""
|
||||
run = self.run
|
||||
auto = self.task.get("auto", {})
|
||||
@@ -583,11 +708,17 @@ class CrawlJob:
|
||||
run["out_dir"] = out_dir
|
||||
os.makedirs(out_dir, exist_ok=True)
|
||||
|
||||
# 恢复持久化状态: 已爬集合 + 上次未爬完的缓存队列
|
||||
visited = set(auto.get("visited", []) or [])
|
||||
pending = auto.get("pending", []) or []
|
||||
# 恢复持久化状态: 已爬集合 + 上次未爬完的缓存队列 (独立状态文件, 不撑大 tasks.json)
|
||||
state = store.load_auto_state(self.task["id"], self.task)
|
||||
visited = set(state.get("visited", []) or [])
|
||||
pending = state.get("pending", []) or []
|
||||
visited_base = len(visited)
|
||||
# 起始网址每次运行都爬(不做去重), 缓存队列继续消费
|
||||
queue = [(seed, 0, "")]
|
||||
# 继续爬取模式(skip_seed): 有缓存队列时跳过起始网址, 直接从待爬队列接着爬
|
||||
skip_seed = bool(self.run.get("skip_seed"))
|
||||
queue = []
|
||||
if not (skip_seed and pending):
|
||||
queue.append((seed, 0, ""))
|
||||
if pending:
|
||||
queue.extend((p["url"], p.get("depth", 0), p.get("source", "")) for p in pending)
|
||||
queued = set(visited)
|
||||
@@ -595,50 +726,70 @@ class CrawlJob:
|
||||
queued.add(normalize_url(u))
|
||||
|
||||
run["progress"]["total"] = len(queue)
|
||||
self._auto_visited = len(visited)
|
||||
self._auto_pending = len(queue)
|
||||
self._persist()
|
||||
|
||||
p, browser, ctx, page, cookie_file = self._open_browser()
|
||||
self._open_browser()
|
||||
try:
|
||||
while queue and not self._stop.is_set():
|
||||
self._wait_if_paused()
|
||||
if max_pages > 0 and len(visited) >= max_pages:
|
||||
self._log("info", f"达到最大页数上限 {max_pages}, 停止")
|
||||
break
|
||||
if max_pages > 0:
|
||||
if skip_seed:
|
||||
# 继续爬取模式: 页数上限按本次新增页数重新计算 (累计 visited 不阻塞继续)
|
||||
if len(visited) - visited_base >= max_pages:
|
||||
self._log("info", f"本次继续爬取达到页数上限 {max_pages}, 停止")
|
||||
break
|
||||
elif len(visited) >= max_pages:
|
||||
self._log("info", f"达到最大页数上限 {max_pages}, 停止")
|
||||
break
|
||||
if self._need_browser_restart(len(visited)):
|
||||
self._log("info", f"已爬 {len(visited) - 1} 页, 主动重启浏览器")
|
||||
self._reopen_browser()
|
||||
url, depth, src = queue.pop(0)
|
||||
key = normalize_url(url)
|
||||
if key in visited and url != seed: # 起始网址不去重, 其余已爬跳过
|
||||
continue
|
||||
visited.add(key)
|
||||
self._auto_visited = len(visited)
|
||||
self._auto_pending = len(queue)
|
||||
idx = len(visited)
|
||||
run["progress"]["current_url"] = url
|
||||
run["progress"]["done"] = len(visited)
|
||||
self._persist()
|
||||
entry = self._retry_crawl(page, ctx, url, idx, out_dir,
|
||||
entry = self._retry_crawl(self._page(), self._ctx(), url, idx, out_dir,
|
||||
source_url=src, depth=depth)
|
||||
run["results"].append(entry)
|
||||
self._bump_stats(entry)
|
||||
self._persist()
|
||||
# 无深度限制或未达深度限制时持续发现链接
|
||||
if entry["status"] == "OK" and (max_depth == 0 or depth < max_depth):
|
||||
for link in self._discover_links(page):
|
||||
for link in self._discover_links(self._page()):
|
||||
lk = normalize_url(link)
|
||||
if lk not in visited and lk not in queued:
|
||||
queued.add(lk)
|
||||
queue.append((link, depth + 1, url))
|
||||
self._auto_pending = len(queue) # 新链接入队后实时刷新
|
||||
if entry["status"] == "OK":
|
||||
self._delay()
|
||||
run["progress"]["total"] = len(visited)
|
||||
if not self._stop.is_set() and len(queue) == 0:
|
||||
self._log("info", f"无新链接可爬, 任务结束 (共 {len(visited)} 页)")
|
||||
finally:
|
||||
self._close_browser(p, browser, ctx, cookie_file)
|
||||
self._close_browser()
|
||||
# 无论完成/停止/异常, 都保存缓存队列与已爬集合, 便于下次继续
|
||||
# 保存前应用当前排除规则过滤, 避免运行中配置的 exclude 被内存快照覆盖
|
||||
try:
|
||||
auto["pending"] = [
|
||||
{"url": u, "depth": d, "source": s} for u, d, s in queue]
|
||||
auto["visited"] = list(visited)
|
||||
store.upsert_task(self.task)
|
||||
db.upsert_task(self.task)
|
||||
auto_cfg = self.task.get("auto") or {}
|
||||
keep = [(u, d, s) for u, d, s in queue if not url_excluded(u, auto_cfg)]
|
||||
if len(keep) != len(queue):
|
||||
self._log("info", f"保存状态时按排除规则过滤 {len(queue) - len(keep)} 条")
|
||||
store.save_auto_state(self.task["id"], {
|
||||
"pending": [
|
||||
{"url": u, "depth": d, "source": s} for u, d, s in keep],
|
||||
"visited": list(visited),
|
||||
})
|
||||
db.upsert_task_async(self.task)
|
||||
except Exception:
|
||||
pass
|
||||
self._persist()
|
||||
@@ -33,11 +33,24 @@ def notify_email(task, run):
|
||||
if len(results) > 50:
|
||||
lines.append(f" ... 共 {len(results)} 条")
|
||||
body = "\n".join(lines)
|
||||
return send_attachment(f"[爬虫完成] {task['name']}", body, to)
|
||||
|
||||
|
||||
def send_attachment(subject, body, to, attach_path=None):
|
||||
"""发送带附件的邮件 (复用 send_email.py), 返回 (bool, msg)
|
||||
attach_path: 附件文件路径, 可多个
|
||||
"""
|
||||
if not os.path.exists(SEND_EMAIL):
|
||||
return False, "send_email.py 不存在"
|
||||
cmd = [sys.executable, SEND_EMAIL, subject, body, "--to", to]
|
||||
if attach_path:
|
||||
if isinstance(attach_path, str):
|
||||
attach_path = [attach_path]
|
||||
for p in attach_path:
|
||||
if os.path.isfile(p):
|
||||
cmd += ["--attach", p]
|
||||
try:
|
||||
r = subprocess.run(
|
||||
[sys.executable, SEND_EMAIL, f"[爬虫完成] {task['name']}", body, "--to", to],
|
||||
timeout=60, capture_output=True,
|
||||
)
|
||||
r = subprocess.run(cmd, timeout=600, capture_output=True)
|
||||
return r.returncode == 0, (r.stderr or b"").decode(errors="ignore")[:200]
|
||||
except Exception as e:
|
||||
return False, str(e)
|
||||
+228
-29
@@ -3,12 +3,13 @@
|
||||
|
||||
const $ = (id) => document.getElementById(id);
|
||||
const MODE_LABEL = { batch: "批量", scheduled: "定时", auto: "自动" };
|
||||
const DETAIL_PAGE_SIZE = 100; // 详情结果每页条数 (默认)
|
||||
|
||||
const state = {
|
||||
tasks: [],
|
||||
editTask: null, // 正在编辑的任务
|
||||
mode: "batch", // 当前表单模式
|
||||
detail: { task: null, runId: null, logOffset: 0, timer: null },
|
||||
detail: { task: null, runId: null, logOffset: 0, timer: null, pageSize: DETAIL_PAGE_SIZE },
|
||||
logTimer: null,
|
||||
};
|
||||
let formDirty = false; // 新建/编辑表单是否有未保存修改
|
||||
@@ -45,26 +46,32 @@ async function loadTasks() {
|
||||
pill.textContent = `运行中任务: ${running}`;
|
||||
pill.classList.toggle("active", running > 0);
|
||||
const st = await api("/api/status");
|
||||
$("version").textContent = `v${st.version}`;
|
||||
$("version").textContent = st && st.version ? `v${st.version}` : "v-";
|
||||
} catch (e) {
|
||||
toast("加载任务失败: " + e.message, true);
|
||||
}
|
||||
loadStats();
|
||||
}
|
||||
|
||||
const STAT_IDS = ["stTasks", "stRunning", "stRuns", "stOk", "stFail", "stImgs", "stDisk"];
|
||||
|
||||
async function loadStats() {
|
||||
// 未读到真实数据前, 一律显示 "-"
|
||||
STAT_IDS.forEach((id) => { $(id).textContent = "-"; });
|
||||
$("trashCount").textContent = "-";
|
||||
try {
|
||||
const s = await api("/api/stats");
|
||||
$("stTasks").textContent = s.tasks;
|
||||
$("stRunning").textContent = s.running;
|
||||
$("stRuns").textContent = s.runs;
|
||||
$("stOk").textContent = s.ok;
|
||||
$("stFail").textContent = s.fail;
|
||||
$("stImgs").textContent = s.images;
|
||||
$("stDisk").textContent = s.disk_mb >= 1024
|
||||
? (s.disk_mb / 1024).toFixed(1) + " GB" : s.disk_mb + " MB";
|
||||
$("trashCount").textContent = s.trash || 0;
|
||||
} catch (e) { /* 统计失败忽略 */ }
|
||||
$("stTasks").textContent = s.tasks ?? "-";
|
||||
$("stRunning").textContent = s.running ?? "-";
|
||||
$("stRuns").textContent = s.runs ?? "-";
|
||||
$("stOk").textContent = s.ok ?? "-";
|
||||
$("stFail").textContent = s.fail ?? "-";
|
||||
$("stImgs").textContent = s.images ?? "-";
|
||||
const mb = Number(s.disk_mb);
|
||||
$("stDisk").textContent = Number.isFinite(mb) && s.disk_mb != null
|
||||
? (mb >= 1024 ? (mb / 1024).toFixed(1) + " GB" : mb + " MB") : "-";
|
||||
$("trashCount").textContent = s.trash ?? "-";
|
||||
} catch (e) { /* 统计失败保持 "-" */ }
|
||||
}
|
||||
|
||||
function taskStatusBadge(t) {
|
||||
@@ -516,13 +523,15 @@ async function submitForm(e) {
|
||||
}
|
||||
|
||||
/* ---------------- 详情 ---------------- */
|
||||
|
||||
async function openDetail(tid) {
|
||||
try {
|
||||
const t = await api(`/api/tasks/${tid}`);
|
||||
const t = await api(`/api/tasks/${tid}?page_size=${DETAIL_PAGE_SIZE}`);
|
||||
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);
|
||||
const cur = t.run;
|
||||
state.detail.runId = cur ? cur.id : (runs[0] ? runs[0].id : null);
|
||||
state.detail.pageSize = (t.run && t.run.run_page_size) || DETAIL_PAGE_SIZE;
|
||||
state.detail.logOffset = 0;
|
||||
$("detailTitle").textContent = `任务详情 · ${t.name}`;
|
||||
renderDetail();
|
||||
@@ -531,11 +540,109 @@ async function openDetail(tid) {
|
||||
} catch (e) { toast(e.message, true); }
|
||||
}
|
||||
|
||||
async function loadRunPage(tid, rid, page, size) {
|
||||
const ps = size || state.detail.pageSize || DETAIL_PAGE_SIZE;
|
||||
return api(`/api/tasks/${tid}?run=${rid}&page=${page}&page_size=${ps}`);
|
||||
}
|
||||
|
||||
async function selectRun(rid) {
|
||||
const t = state.detail.task;
|
||||
if (!t) return;
|
||||
try {
|
||||
const d = await loadRunPage(t.id, rid, 1);
|
||||
state.detail.task = d;
|
||||
state.detail.runId = d.run ? d.run.id : null;
|
||||
state.detail.pageSize = (d.run && d.run.run_page_size) || DETAIL_PAGE_SIZE;
|
||||
state.detail.logOffset = 0;
|
||||
renderDetail();
|
||||
startLogPoll();
|
||||
} catch (e) { toast(e.message, true); }
|
||||
}
|
||||
|
||||
async function goRunPage(p) {
|
||||
const t = state.detail.task;
|
||||
if (!t || !state.detail.runId) return;
|
||||
const ps = (t.run && t.run.run_page_size) || state.detail.pageSize || DETAIL_PAGE_SIZE;
|
||||
try {
|
||||
const d = await loadRunPage(t.id, state.detail.runId, p, ps);
|
||||
state.detail.task = d;
|
||||
state.detail.pageSize = (d.run && d.run.run_page_size) || ps;
|
||||
renderDetail();
|
||||
} catch (e) { toast(e.message, true); }
|
||||
}
|
||||
|
||||
/* 跳转到指定页码 (输入框回车/点击 GO) */
|
||||
function jumpRunPage() {
|
||||
const t = state.detail.task;
|
||||
if (!t || !state.detail.runId || !t.run) return;
|
||||
const inp = $("pagerJump");
|
||||
const pages = t.run.run_pages || 1;
|
||||
let p = parseInt(inp.value, 10);
|
||||
if (!p || isNaN(p)) p = 1;
|
||||
p = Math.min(Math.max(p, 1), pages);
|
||||
inp.value = p;
|
||||
goRunPage(p);
|
||||
}
|
||||
|
||||
/* 切换每页条数 */
|
||||
async function changePageSize(size) {
|
||||
const t = state.detail.task;
|
||||
if (!t || !state.detail.runId) return;
|
||||
try {
|
||||
const d = await loadRunPage(t.id, state.detail.runId, 1, parseInt(size, 10) || 100);
|
||||
state.detail.task = d;
|
||||
state.detail.pageSize = (d.run && d.run.run_page_size) || 100;
|
||||
renderDetail();
|
||||
} catch (e) { toast(e.message, true); }
|
||||
}
|
||||
|
||||
/* 生成页码窗口列表, 0 表示省略号: 1 2 3 … 97 98 99 100 */
|
||||
function pageWindow(cur, pages, width) {
|
||||
width = width || 5;
|
||||
const win = [];
|
||||
if (pages <= width + 2) {
|
||||
for (let i = 1; i <= pages; i++) win.push(i);
|
||||
return win;
|
||||
}
|
||||
win.push(1);
|
||||
let lo = Math.max(2, cur - Math.floor(width / 2));
|
||||
let hi = Math.min(pages - 1, cur + Math.floor(width / 2));
|
||||
if (hi - lo < width - 1) {
|
||||
if (lo <= 2) hi = lo + width - 1;
|
||||
else lo = hi - width + 1;
|
||||
}
|
||||
if (lo > 2) win.push(0);
|
||||
for (let i = lo; i <= hi; i++) win.push(i);
|
||||
if (hi < pages - 1) win.push(0);
|
||||
win.push(pages);
|
||||
return win;
|
||||
}
|
||||
|
||||
/* 重爬失败页: 注入待爬队列并自动继续爬取 */
|
||||
async function retryFailed() {
|
||||
const t = state.detail.task;
|
||||
if (!t || !state.detail.runId) return;
|
||||
if (!confirm("将把本次运行失败的页面重新加入待爬队列并立即继续爬取,确定?")) return;
|
||||
try {
|
||||
const d = await api(`/api/tasks/${t.id}/retry-failed`, {
|
||||
method: "POST",
|
||||
headers: { "Content-Type": "application/json" },
|
||||
body: JSON.stringify({ run: state.detail.runId }),
|
||||
});
|
||||
toast(`已注入 ${d.injected} 个失败页,开始继续爬取…`);
|
||||
if ((d.injected || 0) > 0) {
|
||||
await api(`/api/tasks/${t.id}/continue`, { method: "POST" });
|
||||
}
|
||||
loadTasks();
|
||||
openDetail(t.id);
|
||||
} 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 run = t.run || null;
|
||||
state.detail.runId = run ? run.id : (runs[0] ? runs[0].id : null);
|
||||
const cfg = t.config || {};
|
||||
|
||||
const meta = `
|
||||
@@ -546,8 +653,8 @@ function renderDetail() {
|
||||
<span>输出: <b>${esc(cfg.out_dir || "out/" + t.id)}</b></span>
|
||||
${t.mode === "auto" ? `
|
||||
<span>起始: <b>${esc((t.auto && t.auto.seed_url) || "")}</b></span>
|
||||
<span>待爬缓存: <b>${(t.auto && (t.auto.pending || []).length) || 0}</b> 条(停止后可继续爬取)</span>
|
||||
<span>已爬: <b>${(t.auto && (t.auto.visited || []).length) || 0}</b> 页</span>` : ""}
|
||||
<span>待爬缓存: <b>${t.auto_pending_count ?? "-"}</b> 条(停止后可继续爬取)</span>
|
||||
<span>已爬: <b>${t.auto_visited_count ?? "-"}</b> 页</span>` : ""}
|
||||
</div>`;
|
||||
|
||||
let sched = "";
|
||||
@@ -567,10 +674,15 @@ function renderDetail() {
|
||||
${t.running ? `
|
||||
<button class="btn" onclick="actTask('${t.id}','pause')">⏸ 暂停</button>
|
||||
<button class="btn danger" onclick="actTask('${t.id}','stop')">⏹ 终止</button>`
|
||||
: `<button class="btn primary" onclick="actTask('${t.id}','start')">▶ 立即执行</button>`}
|
||||
: `
|
||||
<button class="btn primary" onclick="actTask('${t.id}','start')">▶ 立即执行</button>
|
||||
${t.mode === "auto" && (t.auto_pending_count || 0) > 0 ? `
|
||||
<button class="btn" onclick="contTask('${t.id}')">⏩ 继续爬取(缓存 ${t.auto_pending_count} 条)</button>` : ""}`}
|
||||
${t.mode === "auto" ? `<button class="btn" onclick="probeTask('${t.id}')">🧪 试爬取</button>
|
||||
<button class="btn" ${t.running ? "disabled" : ""} onclick="clearCache('${t.id}')">🧹 清空缓存</button>` : ""}
|
||||
<button class="btn" onclick="openEdit('${t.id}')">✏️ 编辑配置</button>
|
||||
<button class="btn" onclick="exportZip('${t.id}')">📦 打包下载</button>
|
||||
<button class="btn" onclick="exportEmail('${t.id}')">📧 发邮箱</button>
|
||||
</div>`;
|
||||
|
||||
const chips = runs.map((r) => `
|
||||
@@ -604,7 +716,7 @@ function renderRunPanel(t, run) {
|
||||
const meta = r.meta_file
|
||||
? `<span class="file-link" onclick="showMeta('${t.id}','${esc(r.meta_file)}')">📋 元数据</span>` : "—";
|
||||
return `<tr>
|
||||
<td>${i + 1}</td>
|
||||
<td>${(run.results_total ? (run.run_page - 1) * DETAIL_PAGE_SIZE : 0) + i + 1}</td>
|
||||
<td class="${statusCls}">${r.status}</td>
|
||||
<td class="title-cell" title="${esc(r.title)}">${esc(r.title)}</td>
|
||||
<td class="url-cell" title="${esc(r.url)}">${esc(r.url)}</td>
|
||||
@@ -621,6 +733,35 @@ function renderRunPanel(t, run) {
|
||||
onclick="previewFile('${t.id}','${esc(im.file)}','图片预览')">`).join("")}
|
||||
</div>` : "";
|
||||
|
||||
const total = run.results_total || results.length;
|
||||
const pages = run.run_pages || 1;
|
||||
const page = run.run_page || 1;
|
||||
const pageSize = run.run_page_size || DETAIL_PAGE_SIZE;
|
||||
const multi = pages > 1;
|
||||
const pager = total > 0 ? `
|
||||
<div class="pager">
|
||||
${multi ? `
|
||||
<button class="btn sm" title="首页" ${page <= 1 ? "disabled" : ""} onclick="goRunPage(1)">⏮</button>
|
||||
<button class="btn sm" title="上一页" ${page <= 1 ? "disabled" : ""} onclick="goRunPage(${page - 1})">◀</button>
|
||||
<span class="pager-pages">${pageWindow(page, pages).map((p) => p === 0
|
||||
? '<span class="pager-ellipsis">…</span>'
|
||||
: `<button class="btn sm page-btn ${p === page ? "active" : ""}" onclick="goRunPage(${p})">${p}</button>`).join("")}</span>
|
||||
<button class="btn sm" title="下一页" ${page >= pages ? "disabled" : ""} onclick="goRunPage(${page + 1})">▶</button>
|
||||
<button class="btn sm" title="末页" ${page >= pages ? "disabled" : ""} onclick="goRunPage(${pages})">⏭</button>
|
||||
<span class="pager-info">第 <b>${page}</b> / ${pages} 页 · 共 ${total} 条</span>
|
||||
<span class="pager-jump">跳至
|
||||
<input type="number" id="pagerJump" min="1" max="${pages}" value="${page}"
|
||||
onkeydown="if(event.key==='Enter')jumpRunPage()"> 页
|
||||
<button class="btn sm" onclick="jumpRunPage()">GO</button>
|
||||
</span>`
|
||||
: `<span class="pager-info">共 ${total} 条</span>`}
|
||||
<span class="pager-size">每页
|
||||
<select id="pagerSize" onchange="changePageSize(this.value)">
|
||||
${[50, 100, 200, 500].map((s) => `<option value="${s}" ${s === pageSize ? "selected" : ""}>${s}</option>`).join("")}
|
||||
</select> 条
|
||||
</span>
|
||||
</div>` : "";
|
||||
|
||||
return `
|
||||
<div class="run-panel">
|
||||
<div class="run-stats">
|
||||
@@ -631,6 +772,8 @@ function renderRunPanel(t, run) {
|
||||
<span class="ok">✅ ${s.ok}</span>
|
||||
<span class="fail">❌ ${s.fail}</span>
|
||||
<span class="img">🖼️ ${s.images}</span>
|
||||
${t.mode === "auto" && (s.fail || 0) > 0 && run.status !== "running" && run.status !== "paused" ? `
|
||||
<button class="btn sm" title="将本次运行失败的页面重新加入待爬队列并继续爬取" onclick="retryFailed()">🔄 重爬失败页 (${s.fail})</button>` : ""}
|
||||
</div>
|
||||
${prog}
|
||||
<div class="card-line">当前: <b>${esc(run.progress.current_url || "")}</b></div>
|
||||
@@ -640,7 +783,8 @@ function renderRunPanel(t, run) {
|
||||
<thead><tr><th>#</th><th>状态</th><th>标题</th><th>网址</th><th>爬取时间</th><th>HTML</th><th>TXT</th><th>图片</th><th>元数据</th><th>错误</th></tr></thead>
|
||||
<tbody>${rows}</tbody>
|
||||
</table>
|
||||
</div>` : '<div class="card-line">暂无结果</div>'}
|
||||
</div>
|
||||
${pager}` : '<div class="card-line">暂无结果</div>'}
|
||||
${thumbHtml}
|
||||
<div class="logs" id="logBox"></div>
|
||||
</div>`;
|
||||
@@ -651,13 +795,6 @@ function scrollThumbs() {
|
||||
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;
|
||||
@@ -691,6 +828,68 @@ function stopLogPoll() {
|
||||
if (state.logTimer) { clearInterval(state.logTimer); state.logTimer = null; }
|
||||
}
|
||||
|
||||
async function contTask(tid) {
|
||||
try {
|
||||
const r = await api(`/api/tasks/${tid}/continue`, { method: "POST" });
|
||||
toast("已开始继续爬取缓存队列");
|
||||
loadTasks();
|
||||
if (state.detail.task && state.detail.task.id === tid) openDetail(tid);
|
||||
} catch (e) { toast(e.message, true); }
|
||||
}
|
||||
|
||||
/* ---------------- 打包导出 ---------------- */
|
||||
async function exportZip(tid) {
|
||||
const btn = window.event && window.event.target;
|
||||
if (btn) { btn.disabled = true; btn.textContent = "⏳ 打包中..."; }
|
||||
try {
|
||||
const res = await fetch(`/api/export?task_id=${tid}`);
|
||||
if (!res.ok) {
|
||||
let d = null;
|
||||
try { d = await res.json(); } catch (e) { /* ignore */ }
|
||||
throw new Error((d && d.error) || `HTTP ${res.status}`);
|
||||
}
|
||||
let fname = "export.zip";
|
||||
const disp = res.headers.get("Content-Disposition") || "";
|
||||
const m = disp.match(/filename\*?=(?:UTF-8'')?"?([^";]+)"?/i);
|
||||
if (m) { try { fname = decodeURIComponent(m[1]); } catch (e) { fname = m[1]; } }
|
||||
const blob = await res.blob();
|
||||
const a = document.createElement("a");
|
||||
a.href = URL.createObjectURL(blob);
|
||||
a.download = fname;
|
||||
document.body.appendChild(a);
|
||||
a.click();
|
||||
a.remove();
|
||||
setTimeout(() => URL.revokeObjectURL(a.href), 5000);
|
||||
toast(`打包下载完成 ${fname} (${(blob.size / 1048576).toFixed(1)} MB)`);
|
||||
} catch (e) {
|
||||
toast("打包失败: " + e.message, true);
|
||||
} finally {
|
||||
if (btn) { btn.disabled = false; btn.textContent = "📦 打包下载"; }
|
||||
}
|
||||
}
|
||||
|
||||
async function exportEmail(tid) {
|
||||
const t = (state.detail.task && state.detail.task.id === tid)
|
||||
? state.detail.task : state.tasks.find((x) => x.id === tid);
|
||||
const def = (t && t.config && t.config.notify_email) || "wlq@tphai.com";
|
||||
const email = prompt("发送到邮箱(留空默认 wlq@tphai.com):", def);
|
||||
if (email === null) return; // 用户取消
|
||||
const btn = window.event && window.event.target;
|
||||
if (btn) { btn.disabled = true; btn.textContent = "⏳ 打包发送中..."; }
|
||||
try {
|
||||
const r = await api("/api/export/email", {
|
||||
method: "POST",
|
||||
headers: { "Content-Type": "application/json" },
|
||||
body: JSON.stringify({ task_id: tid, email: email.trim() || "wlq@tphai.com" }),
|
||||
});
|
||||
toast(r.msg || "已发送");
|
||||
} catch (e) {
|
||||
toast("发送失败: " + e.message, true);
|
||||
} finally {
|
||||
if (btn) { btn.disabled = false; btn.textContent = "📧 发邮箱"; }
|
||||
}
|
||||
}
|
||||
|
||||
/* ---------------- 试爬取 ---------------- */
|
||||
function collectProbeFromForm() {
|
||||
const f = $("taskForm");
|
||||
|
||||
+10
-10
@@ -14,9 +14,9 @@
|
||||
<input id="searchInput" placeholder="搜索网址 / 标题 / 任务名..." maxlength="100">
|
||||
<button id="btnSearch" class="btn sm primary" title="搜索">🔍</button>
|
||||
</div>
|
||||
<span class="status-pill" id="statusPill">运行中任务: 0</span>
|
||||
<span class="status-pill" id="statusPill">运行中任务: -</span>
|
||||
<button id="btnTheme" class="btn ghost" title="切换日间/夜间主题">☀️</button>
|
||||
<button id="btnTrash" class="btn ghost" title="回收站">🗑️ 回收站 <span id="trashCount" class="trash-count">0</span></button>
|
||||
<button id="btnTrash" class="btn ghost" title="回收站">🗑️ 回收站 <span id="trashCount" class="trash-count">-</span></button>
|
||||
<button id="btnRefresh" class="btn ghost">⟳ 刷新</button>
|
||||
<button id="btnNew" class="btn primary">+ 新建任务</button>
|
||||
</div>
|
||||
@@ -24,13 +24,13 @@
|
||||
|
||||
<main>
|
||||
<div class="stats-bar" id="statsBar">
|
||||
<div class="stat-card"><div class="stat-num" id="stTasks">0</div><div class="stat-label">📋 任务总数</div></div>
|
||||
<div class="stat-card"><div class="stat-num num-run" id="stRunning">0</div><div class="stat-label">🔄 正在运行</div></div>
|
||||
<div class="stat-card"><div class="stat-num" id="stRuns">0</div><div class="stat-label">📊 累计运行次数</div></div>
|
||||
<div class="stat-card"><div class="stat-num num-ok" id="stOk">0</div><div class="stat-label">✅ 成功页面</div></div>
|
||||
<div class="stat-card"><div class="stat-num num-fail" id="stFail">0</div><div class="stat-label">❌ 失败页面</div></div>
|
||||
<div class="stat-card"><div class="stat-num" id="stImgs">0</div><div class="stat-label">🖼️ 已爬图片</div></div>
|
||||
<div class="stat-card"><div class="stat-num" id="stDisk">0</div><div class="stat-label">💾 磁盘占用</div></div>
|
||||
<div class="stat-card"><div class="stat-num" id="stTasks">-</div><div class="stat-label">📋 任务总数</div></div>
|
||||
<div class="stat-card"><div class="stat-num num-run" id="stRunning">-</div><div class="stat-label">🔄 正在运行</div></div>
|
||||
<div class="stat-card"><div class="stat-num" id="stRuns">-</div><div class="stat-label">📊 累计运行次数</div></div>
|
||||
<div class="stat-card"><div class="stat-num num-ok" id="stOk">-</div><div class="stat-label">✅ 成功页面</div></div>
|
||||
<div class="stat-card"><div class="stat-num num-fail" id="stFail">-</div><div class="stat-label">❌ 失败页面</div></div>
|
||||
<div class="stat-card"><div class="stat-num" id="stImgs">-</div><div class="stat-label">🖼️ 已爬图片</div></div>
|
||||
<div class="stat-card"><div class="stat-num" id="stDisk">-</div><div class="stat-label">💾 磁盘占用</div></div>
|
||||
</div>
|
||||
<div id="taskList" class="task-grid"></div>
|
||||
<div id="emptyState" class="empty hidden">
|
||||
@@ -126,7 +126,7 @@
|
||||
<div class="field"><label>包含规则(每行一个,子串或正则)</label>
|
||||
<textarea name="include" rows="3" placeholder="techpowerup.com/review /news/"></textarea></div>
|
||||
<div class="field"><label>排除规则</label>
|
||||
<textarea name="exclude" rows="3" placeholder="login, signup, /tag/, /forum/"></textarea></div>
|
||||
<textarea name="exclude" rows="3" placeholder="每行一个关键词,URL 含任一关键词即不爬取 如: MyComments.html、OtherPosts.html、/comments 勾选正则后按正则匹配"></textarea></div>
|
||||
</div>
|
||||
<div class="row2">
|
||||
<div class="field check"><label><input name="same_domain" type="checkbox" checked> 仅爬同域名</label></div>
|
||||
|
||||
@@ -275,3 +275,32 @@ td.title-cell { max-width: 220px; overflow: hidden; text-overflow: ellipsis; }
|
||||
@keyframes fadein { from { opacity: 0; transform: translate(-50%, -8px); } }
|
||||
::-webkit-scrollbar { width: 8px; height: 8px; }
|
||||
::-webkit-scrollbar-thumb { background: var(--thumb); border-radius: 4px; }
|
||||
|
||||
/* 结果分页控件 */
|
||||
.pager {
|
||||
display: flex;
|
||||
align-items: center;
|
||||
justify-content: center;
|
||||
gap: 12px;
|
||||
padding: 10px 0;
|
||||
flex-wrap: wrap;
|
||||
}
|
||||
.pager-info { color: var(--text-dim, #999); font-size: 13px; }
|
||||
.pager .btn:disabled { opacity: 0.4; cursor: not-allowed; }
|
||||
.pager-pages { display: flex; align-items: center; gap: 4px; flex-wrap: wrap; }
|
||||
.pager-pages .page-btn { min-width: 30px; padding: 4px 6px; }
|
||||
.pager-pages .page-btn.active {
|
||||
background: var(--accent, #2d6cdf);
|
||||
color: #fff;
|
||||
border-color: var(--accent, #2d6cdf);
|
||||
font-weight: 600;
|
||||
}
|
||||
.pager-ellipsis { color: var(--text-dim, #999); padding: 0 2px; user-select: none; }
|
||||
.pager-jump { display: inline-flex; align-items: center; gap: 4px; color: var(--text-dim, #999); font-size: 13px; }
|
||||
.pager-jump input {
|
||||
width: 56px;
|
||||
padding: 3px 6px;
|
||||
text-align: center;
|
||||
}
|
||||
.pager-size { display: inline-flex; align-items: center; gap: 4px; color: var(--text-dim, #999); font-size: 13px; }
|
||||
.pager-size select { padding: 3px 4px; width: auto; }
|
||||
@@ -10,6 +10,7 @@ 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, ...]}
|
||||
AUTO_STATE_DIR = os.path.join(DATA_DIR, "auto_state") # auto 任务待爬/已爬队列独立存储
|
||||
|
||||
_lock = threading.RLock()
|
||||
|
||||
@@ -140,6 +141,46 @@ def purge_trash():
|
||||
return purged
|
||||
|
||||
|
||||
# ---------------- auto 任务状态 (待爬队列/已爬集合) ----------------
|
||||
# 独立于 tasks.json 存储: 避免上万条 URL 撑大任务文件导致全量读写变慢
|
||||
|
||||
def auto_state_path(task_id):
|
||||
return os.path.join(AUTO_STATE_DIR, f"{task_id}.json")
|
||||
|
||||
|
||||
def load_auto_state(task_id, fallback_task=None):
|
||||
"""读取 auto 任务的 pending/visited 状态; 无独立文件时从旧任务对象迁移"""
|
||||
with _lock:
|
||||
path = auto_state_path(task_id)
|
||||
if os.path.exists(path):
|
||||
st = _load(path, {}) or {}
|
||||
return {"pending": st.get("pending", []), "visited": st.get("visited", [])}
|
||||
if fallback_task:
|
||||
auto = fallback_task.get("auto") or {}
|
||||
st = {"pending": auto.get("pending", []), "visited": auto.get("visited", [])}
|
||||
if st["pending"] or st["visited"]:
|
||||
_save(path, st) # 自动迁移旧数据
|
||||
return st
|
||||
return {"pending": [], "visited": []}
|
||||
|
||||
|
||||
def save_auto_state(task_id, state):
|
||||
with _lock:
|
||||
_save(auto_state_path(task_id), {
|
||||
"pending": state.get("pending", []),
|
||||
"visited": state.get("visited", []),
|
||||
})
|
||||
|
||||
|
||||
def clear_auto_state(task_id):
|
||||
"""删除 auto 状态文件"""
|
||||
with _lock:
|
||||
try:
|
||||
os.remove(auto_state_path(task_id))
|
||||
except OSError:
|
||||
pass
|
||||
|
||||
|
||||
# ---------------- 运行记录 ----------------
|
||||
|
||||
def load_runs_map():
|
||||
|
||||
Reference in New Issue
Block a user