Compare commits

..
11 Commits
Author SHA1 Message Date
hz4th_coder 32eebf3dd1 auto 任务运行中实时显示已爬/待爬缓存数; 人工停止后继续爬取验证
- engine: CrawlJob 增加 _auto_visited/_auto_pending 运行中实时状态 (含发现新链接入队后刷新)
- app.py 详情接口: 任务运行中优先读内存实时值, 不再显示上次运行结束的旧数字
- 端到端验证: 人工停止(已爬3,缓存18) → 继续爬取(跳过seed,爬完18页,缓存0) 
- 运行中采样: visited 1→2 实时增长, pending 15→14 同步递减 
2026-08-12 10:19:19 +08:00
hz4th_coder d9c9f0c633 chore: data/auto_state 为运行时爬取状态, 移出版本控制 (.gitignore) 2026-08-12 09:57:53 +08:00
hz4th_coder 909e8e01b5 新增继续爬取: auto 任务达 max_pages 上限停止后, 有待爬缓存链接时可在详情页一键继续
- POST /api/tasks/<tid>/continue: 校验 auto 模式 + pending 非空, skip_seed 启动
- engine skip_seed: 跳过起始网址直接从缓存队列消费 (不重爬 seed, 不重复入 visited)
- 关键修复: max_pages 原为累计上限(visited>=max_pages), 已达标任务无法继续; 继续模式改为按本次新增页数重新计算上限, 可反复继续直到缓存耗尽
- 详情弹窗 auto 任务且有缓存时显示 [ 继续爬取(缓存 N 条)] 按钮
- 端到端测试: max_pages=2 链路 执行(seed+p1,缓存4) → 继续(p2,p3,缓存2) → 继续(p4,p5,缓存0) → 无缓存报错 
2026-08-12 09:57:46 +08:00
hz4th_coder e749d70a43 chore: data/exports 为运行时打包产物, 移出版本控制 (.gitignore) 2026-08-12 09:43:32 +08:00
hz4th_coder 56f30f78c4 新增打包导出: 任务输出目录打包为 zip, 支持网页下载或发送到指定邮箱(默认 wlq@tphai.com)
- GET /api/export?task_id= 打包输出目录为 zip 并下载 (zip 内按任务名建顶层目录)
- POST /api/export/email 打包后通过 send_email.py 发送附件到邮箱, 邮箱可指定, 默认 wlq@tphai.com
- 详情弹窗新增 [📦 打包下载] [📧 发邮箱] 按钮, 打包中按钮禁用防重复点击
- notify.py 重构: send_attachment(subject, body, to, attach) 通用附件发送
- 打包文件存 data/exports/, 自动清理只保留最近 10 个
- 实测: cnblogs-auto 66MB/1606文件 -> 11.9MB zip 仅1.5s; 邮件发送成功 2.2s
2026-08-12 09:41:59 +08:00
hz4th_coder 0482284cc0 反爬拦截不再重试: 判定为反爬类错误(403/429/503/Cloudflare/captcha/验证页等)直接放弃该页, 节省重试时间; 验证页持续8秒即判定反爬, 不再干等满超时(默认60s)
- engine.is_anti_crawl_error: 反爬错误信号关键词判定
- _retry_crawl: 反爬错误直接 break, entry.error 标注'反爬拦截, 跳过重试'
- _wait_page_settle/_settle_wait: 连续8秒检测到验证页即判定反爬尽早返回
- 集成测试: 反爬页 attempts=1 且 8.1s 放弃; 非反爬错误仍 attempts=retry_count+1 正常重试
2026-08-12 09:37:23 +08:00
hz4th_coder 796e667533 性能优化: 任务列表瘦身+详情分页+auto状态独立存储+磁盘统计缓存+db异步同步; 统计数值未读取时显示-
- /api/tasks 响应 1.57MB -> 12KB (latest_run 只带摘要, 一次读 runs.json)
- /api/stats 磁盘占用加 30s 缓存, 去除重复全量读
- 详情接口分页返回 results (默认100条/页), 前端表格分页+页码跳转
- auto 任务 pending/visited 队列迁移到 data/auto_state/ 独立文件 (tasks.json 1.6MB -> 14KB)
- MySQL 同步改后台线程异步执行 (db.sync_run_async/upsert_task_async), 不再阻塞爬虫和 API
- 启动时自动迁移存量 auto 状态
- 统计/回收站/运行中数值在未读到真实数据前显示 '-'
2026-08-12 09:28:16 +08:00
hz4th_coder 66ba1a7cfc v1.0.14: 自动爬取支持深度/页数无限制(默认0/1000), 待爬链接持久化缓存队列(停止后续爬), 起始网址不去重, 清空缓存API 2026-08-11 21:07:31 +08:00
hz4th_coder 3ad3d00ca8 v1.0.13: 自动爬取改为持续递归(无深度限制, 爬到无新链接为止), max_pages作安全上限(默认200) 2026-08-11 19:45:51 +08:00
hz4th_coder aa0c2f56d0 v1.0.12: crawl_results 合并 html_file/txt_file/meta_file 为单一 base_file(同前缀不同后缀), 含存量数据自动迁移 2026-08-11 19:35:55 +08:00
hz4th_coder bf6be47c01 v1.0.11: 新增历史记录回填脚本 backfill.py(幂等, 支持指定任务), 补齐数据库功能上线前的爬取记录 2026-08-11 19:29:57 +08:00
11 changed files with 746 additions and 107 deletions
+2
View File
@@ -4,3 +4,5 @@ data/*.json
data/cookies_*.json
logs/
out/
data/exports/
data/auto_state/
+4 -2
View File
@@ -49,12 +49,14 @@
- 可启用/停用调度,自动计算下次执行时间;到点自动开跑,跑完自动计算下一次
### 4. 自动爬取模式
给定一个起始网址,系统自动从页面里发现链接、按规则筛选后 BFS 爬取:
给定一个起始网址,系统**持续递归**爬取:爬取页面 → 自动发现符合规则的链接 → 继续爬取 → 继续发现……直到**无新链接可爬**时自动结束。
- **🧪 试爬取**:先用起始网址试跑一次,展示规则筛选后的链接清单(将爬取哪些、排除哪些及原因),确认规则符合预期后再正式开爬
- **最大爬取深度**:0=无限制(默认),N=只爬 N 层
- **最大页数**:0=无限制,默认 1000 作安全上限(防止动态无限链接的站点失控)
- **🔄 缓存续爬**:提取但未爬取的链接自动存入缓存队列(连同已爬集合一并持久化),任务停止/中断后再次运行,从缓存队列**继续爬取**,已爬过的不会重复;起始网址每次运行都重新爬取(不去重);规则变更后可点「🧹 清空缓存」重新开始
- **包含规则**:只爬包含指定子串(或正则)的链接
- **排除规则**:跳过匹配的链接(如 login、/tag/
- **仅同域名**:限制在起始网站内
- **最大页数 / 最大深度**:控制爬取规模
- 其余参数(间隔、重试、图片、通知)同批量模式
### 5. 资源操作信息(元数据)
+293 -36
View File
@@ -4,13 +4,17 @@
启动: /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
import notify
from engine import CrawlJob, probe_links
from scheduler import Scheduler, cron_next, interval_delta
@@ -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"])
@@ -267,7 +355,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 +371,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 +394,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 +407,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 +435,23 @@ def api_trash_clear():
# ---------------- API: 运行控制 ----------------
@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 +492,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 +513,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,
})
@@ -448,6 +563,28 @@ def api_probe():
return jsonify(result)
@app.route("/api/tasks/<tid>/clear_cache", methods=["POST"])
def api_clear_cache(tid):
"""清空自动任务的待爬缓存队列与已爬集合 (规则变更后重新开始用)"""
task = store.get_task(tid)
if not task:
return jsonify({"error": "任务不存在"}), 404
if task.get("mode") != "auto":
return jsonify({"error": "仅自动爬取任务支持清空缓存"}), 400
with JOBS_LOCK:
job = JOBS.get(tid)
if job and job.is_running():
return jsonify({"error": "任务正在运行,无法清空缓存"}), 409
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_async(task)
return jsonify({"ok": True, "msg": "缓存队列已清空"})
@app.route("/api/tasks/<tid>/probe", methods=["POST"])
def api_task_probe(tid):
"""对已保存的自动任务执行试爬取 (使用保存的规则)"""
@@ -479,6 +616,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 # 回收站任务不参与搜索
@@ -488,7 +626,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 ""
@@ -567,18 +705,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)
+26
View File
@@ -0,0 +1,26 @@
#!/usr/bin/env python3
# -*- coding: utf-8 -*-
"""
历史爬取记录回填数据库 (幂等, 可重复执行)
用法:
python backfill.py # 回填所有任务的历史记录
python backfill.py <task_id> ... # 只回填指定任务
说明:
- 将本地 data/ 中的任务/运行/爬取结果(含成功与失败)全量写入 MySQL
- 网页内容不入库, 只写元数据; 重复执行不会产生重复记录
"""
import sys
import db
import store
if __name__ == "__main__":
ids = [a for a in sys.argv[1:] if a.strip()] or None
db.init_db()
stats = db.sync_all_history(ids)
if ids:
print(f"[backfill] 已回填 {len(ids)} 个任务: {stats}")
else:
print(f"[backfill] 已回填全部任务: {stats}")
+92 -7
View File
@@ -5,9 +5,11 @@ MySQL 记录层: 任务 / 运行 / 爬取结果 写入数据库
- 爬取成功/失败均有 status 标记 (OK / FAIL), 运行记录有 run 状态标记
- 所有操作容错: 数据库不可用时不影响爬取主流程 (仅打印日志)
"""
import copy
import json
import os
import time
from concurrent.futures import ThreadPoolExecutor
import pymysql
@@ -67,9 +69,7 @@ DDL = [
source_url VARCHAR(2000),
depth INT,
attempts INT DEFAULT 1,
html_file VARCHAR(500),
txt_file VARCHAR(500),
meta_file VARCHAR(500),
base_file VARCHAR(500),
image_count INT DEFAULT 0,
image_files TEXT,
UNIQUE KEY uk_run_url (run_id, url(500))
@@ -78,6 +78,17 @@ DDL = [
]
def _migrate(conn):
"""存量表结构迁移: 合并 html_file/txt_file/meta_file 为 base_file"""
with conn.cursor() as cur:
cur.execute("SHOW COLUMNS FROM crawl_results LIKE 'html_file'")
if cur.fetchone():
cur.execute("ALTER TABLE crawl_results ADD COLUMN base_file VARCHAR(500) NULL AFTER meta_file")
cur.execute("UPDATE crawl_results SET base_file = REPLACE(html_file, '.html', '') WHERE base_file IS NULL")
cur.execute("ALTER TABLE crawl_results DROP COLUMN html_file, DROP COLUMN txt_file, DROP COLUMN meta_file")
print("[db] 表结构迁移完成: html_file/txt_file/meta_file -> base_file")
def _conn():
cfg = dict(DB_CONFIG)
cfg["database"] = DB_NAME
@@ -92,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:
@@ -100,6 +144,7 @@ def init_db():
cur.execute(f"CREATE DATABASE IF NOT EXISTS `{DB_NAME}` DEFAULT CHARACTER SET utf8mb4")
conn.close()
conn = _conn()
_migrate(conn)
with conn.cursor() as cur:
for ddl in DDL:
cur.execute(ddl)
@@ -219,8 +264,21 @@ def upsert_run(run):
# ---------------- 爬取结果 (每页一条, 成功失败均记录) ----------------
def _base_of(entry):
"""从结果条目提取基础文件名 (html/txt/meta 三个后缀共用同一前缀)"""
for f in (entry.get("meta_file"), entry.get("html_file"), entry.get("txt_file")):
f = f or ""
if f.endswith(".meta.json"):
return f[:-10]
if f.endswith(".html"):
return f[:-5]
if f.endswith(".txt"):
return f[:-4]
return ""
def insert_results(run, results):
"""批量插入爬取结果 (增量)"""
"""批量插入爬取结果 (增量); 文件只记基础名 base_file"""
if not results:
return
@@ -234,7 +292,7 @@ def insert_results(run, results):
r.get("status", "FAIL"), r.get("error"),
_dt(r.get("crawl_time")), (r.get("source_url") or "")[:2000],
r.get("depth"), r.get("attempts", 1),
r.get("html_file", ""), r.get("txt_file", ""), r.get("meta_file", ""),
_base_of(r),
len(r.get("images", []) or []),
_j([im.get("file") for im in (r.get("images") or [])]),
))
@@ -243,8 +301,8 @@ def insert_results(run, results):
"""INSERT IGNORE INTO crawl_results
(run_id, task_id, mode, url, title, status, error,
crawl_time, source_url, depth, attempts,
html_file, txt_file, meta_file, image_count, image_files)
VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s)""",
base_file, image_count, image_files)
VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s)""",
rows,
)
conn.close()
@@ -262,3 +320,30 @@ def sync_run(run, persist):
insert_results(run, results[synced:])
run["_db_count"] = len(results)
persist(run.get("task_id"), run)
# ---------------- 历史数据回填 (幂等) ----------------
def sync_all_history(task_ids=None):
"""把本地 JSON 中的历史任务/运行/结果全量回填数据库
可重复执行 (INSERT IGNORE + 唯一键去重)
task_ids: 指定只回填的任务ID列表, 默认全部
返回统计 dict
"""
import store as _store
stats = {"tasks": 0, "runs": 0, "results": 0}
for task in _store.load_tasks():
if task_ids and task["id"] not in task_ids:
continue
upsert_task(task)
stats["tasks"] += 1
for run in _store.get_runs(task["id"]):
upsert_run(run)
stats["runs"] += 1
results = run.get("results", [])
if results:
insert_results(run, results)
stats["results"] += len(results)
run["_db_count"] = len(results)
_store.save_run(task["id"], run) # 记录已同步标记, 避免运行中重复插入
return stats
+91 -13
View File
@@ -21,6 +21,7 @@ from playwright_stealth import Stealth
import notify
import store
import db
HERE = os.path.dirname(os.path.abspath(__file__))
DATA_DIR = os.path.join(HERE, "data")
@@ -35,6 +36,19 @@ _CHALLENGE_MARKS = [
"checking your browser", "enable javascript and cookies",
]
# 反爬拦截类错误信号: 命中后不再重试 (重试无意义且拖慢任务)
_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)
def is_challenge_page(title, html):
"""判断是否仍在反爬验证页 (仅关键词启发, 避免误伤正常小页面)"""
@@ -55,6 +69,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)
@@ -64,8 +79,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:
@@ -207,6 +226,9 @@ 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
@@ -305,6 +327,7 @@ class CrawlJob:
def _wait_page_settle(self, page, timeout_s):
last_title, stable = "", 0
challenge_hits = 0
start = time.time()
while time.time() - start < timeout_s:
if self._stop.is_set():
@@ -317,8 +340,12 @@ class CrawlJob:
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:
@@ -459,6 +486,11 @@ class CrawlJob:
entry["error"] = str(e)
entry["crawl_time"] = store.now_str()
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()
@@ -567,33 +599,65 @@ class CrawlJob:
return included
def _crawl_auto(self):
"""自动爬取: 持续递归 爬取->发现链接->爬取... 直到无新链接可爬或达到上限
- max_depth: 0=无限制, N=只爬 N 层
- max_pages: 0=无限制, N=安全上限 (继续爬取模式按本次新增页数重新计算)
- 提取但未爬取的链接持久化到 auto.pending, 已爬集合持久化到 auto.visited
- 停止后再次运行从缓存队列继续爬; 起始网址每次运行都重新爬(不去重)
- skip_seed(继续爬取): 跳过起始网址直接消费缓存队列, 页数上限按本次新增重新计算
"""
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
max_pages = int(auto.get("max_pages", 1000) or 0) # 0 = 无限制
max_depth = int(auto.get("max_depth", 0) or 0) # 0 = 无限制
out_dir = self._resolve_out_dir()
run["out_dir"] = out_dir
os.makedirs(out_dir, exist_ok=True)
# 恢复持久化状态: 已爬集合 + 上次未爬完的缓存队列 (独立状态文件, 不撑大 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)
# 起始网址每次运行都爬(不做去重), 缓存队列继续消费
# 继续爬取模式(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)
for u, _d, _s in queue:
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()
queue = [(seed, 0, "")] # (url, depth, 来源链接)
visited = set() # 规范化 URL 去重
queued = set([normalize_url(seed)])
idx = 0
try:
while queue and not self._stop.is_set():
self._wait_if_paused()
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
url, depth, src = queue.pop(0)
key = normalize_url(url)
if key in visited:
if key in visited and url != seed: # 起始网址不去重, 其余已爬跳过
continue
if len(visited) >= max_pages:
break
visited.add(key)
idx += 1
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()
@@ -602,15 +666,29 @@ class CrawlJob:
run["results"].append(entry)
self._bump_stats(entry)
self._persist()
if entry["status"] == "OK" and depth < max_depth:
# 无深度限制或未达深度限制时持续发现链接
if entry["status"] == "OK" and (max_depth == 0 or depth < max_depth):
for link in self._discover_links(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)
run["progress"]["total"] = len(visited)
# 无论完成/停止/异常, 都保存缓存队列与已爬集合, 便于下次继续
try:
store.save_auto_state(self.task["id"], {
"pending": [
{"url": u, "depth": d, "source": s} for u, d, s in queue],
"visited": list(visited),
})
db.upsert_task_async(self.task)
except Exception:
pass
self._persist()
+17 -4
View File
@@ -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)
+154 -33
View File
@@ -45,26 +45,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) {
@@ -82,7 +88,7 @@ function taskStatusBadge(t) {
const GROUPS = [
{ key: "batch", label: "📄 批量爬取", desc: "一次性爬取指定网址列表" },
{ key: "scheduled", label: "⏰ 定时爬取", desc: "按间隔或 cron 表达式定时执行" },
{ key: "auto", label: "🤖 自动爬取", desc: "从起始网址自动发现链接并爬取" },
{ key: "auto", label: "🤖 自动爬取", desc: "从起始网址自动发现链接,持续递归爬取到无新链接为止" },
];
function renderTasks() {
@@ -434,8 +440,8 @@ async function openEdit(tid) {
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;
f.elements["max_pages"].value = t.auto.max_pages ?? 1000;
f.elements["max_depth"].value = t.auto.max_depth ?? 0;
}
syncScheduleUI();
$("formHint").textContent = t.running
@@ -481,8 +487,8 @@ async function submitForm(e) {
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,
max_pages: parseInt(f.elements["max_pages"].value) || 0,
max_depth: parseInt(f.elements["max_depth"].value) || 0,
};
}
if (mode === "scheduled") {
@@ -516,13 +522,15 @@ async function submitForm(e) {
}
/* ---------------- 详情 ---------------- */
const DETAIL_PAGE_SIZE = 100; // 详情结果每页条数
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.logOffset = 0;
$("detailTitle").textContent = `任务详情 · ${t.name}`;
renderDetail();
@@ -531,11 +539,38 @@ async function openDetail(tid) {
} catch (e) { toast(e.message, true); }
}
async function loadRunPage(tid, rid, page) {
return api(`/api/tasks/${tid}?run=${rid}&page=${page}&page_size=${DETAIL_PAGE_SIZE}`);
}
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.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;
try {
const d = await loadRunPage(t.id, state.detail.runId, p);
state.detail.task = d;
renderDetail();
} 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 = `
@@ -544,7 +579,10 @@ function renderDetail() {
<span>任务ID: <b>${t.id}</b></span>
<span>创建: <b>${fmtTime(t.created_at)}</b></span>
<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>` : ""}
${t.mode === "auto" ? `
<span>起始: <b>${esc((t.auto && t.auto.seed_url) || "")}</b></span>
<span>待爬缓存: <b>${t.auto_pending_count ?? "-"}</b> 条(停止后可继续爬取)</span>
<span>已爬: <b>${t.auto_visited_count ?? "-"}</b> 页</span>` : ""}
</div>`;
let sched = "";
@@ -564,9 +602,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>`}
${t.mode === "auto" ? `<button class="btn" onclick="probeTask('${t.id}')">🧪 试爬取</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) => `
@@ -600,7 +644,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>
@@ -617,6 +661,18 @@ 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 pager = total > 0 ? `
<div class="pager">
${total > results.length ? `
<button class="btn sm" ${page <= 1 ? "disabled" : ""} onclick="goRunPage(${page - 1})">◀ 上一页</button>
<span class="pager-info">第 <b>${page}</b> / ${pages} 页 · 共 ${total} 条</span>
<button class="btn sm" ${page >= pages ? "disabled" : ""} onclick="goRunPage(${page + 1})">下一页 ▶</button>`
: `<span class="pager-info">共 ${total} 条</span>`}
</div>` : "";
return `
<div class="run-panel">
<div class="run-stats">
@@ -636,7 +692,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>`;
@@ -647,13 +704,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;
@@ -687,6 +737,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");
@@ -721,6 +833,15 @@ async function probeTask(tid) {
} catch (e) { toast(e.message, true); }
}
async function clearCache(tid) {
if (!confirm("清空待爬缓存与已爬记录?\n下次运行将从起始网址重新开始爬取。")) return;
try {
const r = await api(`/api/tasks/${tid}/clear_cache`, { method: "POST" });
toast(r.msg || "已清空");
openDetail(tid);
} catch (e) { toast(e.message, true); }
}
function renderProbe(r) {
const body = $("probeBody");
if (!r.ok) {
+14 -12
View File
@@ -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">
@@ -116,7 +116,7 @@
</div>
<div id="autoBox" class="hidden box">
<div class="field"><label>起始网址 *系统将自动发现符合规则的链接并爬取</label>
<div class="field"><label>起始网址 *(自动发现符合规则的链接,持续递归爬取直到无新链接可爬</label>
<div class="row2">
<input name="seed_url" placeholder="https://example.com/news" style="flex:1">
<button type="button" class="btn sm" id="btnProbe">🧪 试爬取</button>
@@ -133,8 +133,10 @@
<div class="field check"><label><input name="use_regex" type="checkbox"> 规则按正则匹配</label></div>
</div>
<div class="row2">
<div class="field"><label>最大页数</label><input name="max_pages" type="number" value="50"></div>
<div class="field"><label>最大爬取深度</label><input name="max_depth" type="number" value="2"></div>
<div class="field"><label>最大爬取深度(0=无限制,N=只爬 N 层)</label>
<input name="max_depth" type="number" min="0" value="0"></div>
<div class="field"><label>最大页数(0=无限制,默认 1000 作安全上限)</label>
<input name="max_pages" type="number" min="0" value="1000"></div>
</div>
</div>
</div>
+12
View File
@@ -275,3 +275,15 @@ 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; }
+41
View File
@@ -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():