# -*- coding: utf-8 -*- """ 通用爬虫系统 - Web 管理后端 启动: /home/hz1/miniconda3/envs/openclaw/bin/python app.py (默认端口 16062) """ import os import threading from datetime import datetime from flask import Flask, jsonify, request, send_file, send_from_directory import store import db from engine import CrawlJob, probe_links from scheduler import Scheduler, cron_next, interval_delta HERE = os.path.dirname(os.path.abspath(__file__)) PORT = int(os.environ.get("CRAWLER_PORT", "16062")) app = Flask(__name__, static_folder="static", static_url_path="") JOBS = {} # task_id -> CrawlJob JOBS_LOCK = threading.RLock() STARTED = datetime.now() DEFAULT_CONFIG = { "out_dir": "", # 留空 -> out/<任务ID> "delay_min": 2, "delay_max": 5, "timeout": 60, "crawl_images": False, "retry_count": 2, "retry_interval": 3, "notify": False, "notify_email": "wlq@tphai.com", } def now_str(): return store.now_str() def _parse_first_run(s): """解析表单提交的首次执行时间 (datetime-local 格式), 非法返回 None""" if not s: return None try: return datetime.strptime(str(s), "%Y-%m-%dT%H:%M") except Exception: return None def _schedule_next_run(sch): """根据调度配置 + 首次执行时间计算 next_run (str)""" first_run = _parse_first_run(sch.get("first_run")) base = first_run if (first_run and first_run > datetime.now()) else None if sch.get("type") == "cron": expr = sch.get("cron") or "0 * * * *" nn = cron_next(expr, base or datetime.now()) if not nn: raise ValueError("cron 表达式在未来一年内无匹配时间") return nn.strftime("%Y-%m-%d %H:%M:%S") sch.setdefault("interval_unit", "hours") sch.setdefault("interval_value", 24) if base: return base.strftime("%Y-%m-%d %H:%M:%S") return (datetime.now() + interval_delta(sch)).strftime("%Y-%m-%d %H:%M:%S") def resolve_out_dir(task): cfg = task.get("config", {}) or {} if cfg.get("out_dir", "").strip(): return cfg["out_dir"].strip() return os.path.join(HERE, "out", task["id"]) def make_run(task): return { "id": store.new_id("r"), "task_id": task["id"], "mode": task.get("mode", "batch"), "status": "running", "progress": {"done": 0, "total": 0, "current_url": "", "percent": 0}, "started_at": now_str(), "finished_at": "", "stats": {"ok": 0, "fail": 0, "images": 0}, "results": [], "logs": [], "out_dir": resolve_out_dir(task), } def persist_cb(task_id, run): total = run["progress"].get("total") or 0 done = run["progress"].get("done") or 0 run["progress"]["percent"] = round(done * 100 / total) if total else 0 store.save_run(task_id, run) try: db.sync_run(run, lambda tid, r: store.save_run(tid, r)) except Exception as e: print(f"[db] 同步失败: {e}", flush=True) def start_run(task): """为任务启动一次爬取, 返回 (run, error)""" with JOBS_LOCK: job = JOBS.get(task["id"]) if job and job.is_running(): return None, "该任务已有正在运行的爬取" run = make_run(task) store.add_run(task["id"], run) job = CrawlJob(task, run, persist_cb) JOBS[task["id"]] = job job.start() db.upsert_task(task) # 确保任务在库中 return run, None def _stop_job(tid): with JOBS_LOCK: job = JOBS.get(tid) if job and job.is_running(): job.stop() return True return False # ---------------- 页面 ---------------- @app.route("/") def index(): return send_from_directory(app.static_folder, "index.html") # ---------------- API: 状态 ---------------- @app.route("/api/status") def api_status(): with JOBS_LOCK: running = sum(1 for j in JOBS.values() if j.is_running()) return jsonify({ "app": "universal-crawler", "version": "1.0.0", "running": running, "uptime": str(datetime.now() - STARTED).split(".")[0], }) # ---------------- API: 任务 ---------------- @app.route("/api/tasks", methods=["GET"]) def api_tasks(): tasks = [t for t in store.load_tasks() if not t.get("deleted_at")] for t in tasks: runs = store.get_runs(t["id"]) t["latest_run"] = runs[-1] if runs else None with JOBS_LOCK: job = JOBS.get(t["id"]) t["running"] = bool(job and job.is_running()) tasks.sort(key=lambda x: x.get("created_at", ""), reverse=True) return jsonify(tasks) def _clean_urls(lines): """清洗网址行: 去空白, 去行内 # 注释 (URL 本身不含空格, 安全), 自动补 https 前缀""" out = [] for u in lines or []: u = str(u).split(" #")[0].strip() if not u: continue if not u.startswith("http"): u = "https://" + u out.append(u) return out @app.route("/api/tasks", methods=["POST"]) def api_create_task(): body = request.get_json(force=True) or {} name = (body.get("name") or "").strip() mode = body.get("mode", "batch") if not name: return jsonify({"error": "项目名不能为空"}), 400 if mode not in ("batch", "auto", "scheduled"): return jsonify({"error": "无效的模式"}), 400 task = { "id": store.new_id("t"), "name": name, "mode": mode, "created_at": now_str(), "updated_at": now_str(), "config": {**DEFAULT_CONFIG, **(body.get("config") or {})}, "urls": _clean_urls(body.get("urls") or []), } if mode == "auto": auto = dict(body.get("auto") or {}) seed = (auto.get("seed_url") or "").strip() if not seed: return jsonify({"error": "自动模式需要填写起始网址"}), 400 if not seed.startswith("http"): seed = "https://" + seed auto["seed_url"] = seed task["auto"] = auto elif mode == "scheduled": sch = dict(body.get("schedule") or {}) sch.setdefault("enabled", True) sch.setdefault("type", "interval") try: sch["next_run"] = _schedule_next_run(sch) except ValueError as e: return jsonify({"error": f"调度配置错误: {e}"}), 400 sch.setdefault("last_run", "") sch.setdefault("runs_count", 0) task["schedule"] = sch if not task["urls"]: return jsonify({"error": "定时任务需要网址列表"}), 400 else: if not task["urls"]: return jsonify({"error": "请至少填写一个网址"}), 400 store.upsert_task(task) db.upsert_task(task) return jsonify(task), 201 @app.route("/api/tasks/", methods=["GET"]) def api_task_detail(tid): task = store.get_task(tid) if not task: return jsonify({"error": "任务不存在"}), 404 task["runs"] = list(reversed(store.get_runs(tid))) with JOBS_LOCK: job = JOBS.get(tid) task["running"] = bool(job and job.is_running()) return jsonify(task) @app.route("/api/tasks/", methods=["PUT"]) def api_update_task(tid): task = store.get_task(tid) if not task: return jsonify({"error": "任务不存在"}), 404 if task.get("deleted_at"): return jsonify({"error": "任务在回收站中,请先恢复"}), 400 body = request.get_json(force=True) or {} with JOBS_LOCK: job = JOBS.get(tid) running = bool(job and job.is_running()) if "name" in body and str(body["name"]).strip(): task["name"] = str(body["name"]).strip() if "urls" in body: task["urls"] = _clean_urls(body["urls"]) if "config" in body: merged = {**task.get("config", {}), **body["config"]} task["config"] = merged if running: job.update_config(body["config"]) # 运行中热更新 if "auto" in body and task.get("mode") == "auto": task["auto"] = {**task.get("auto", {}), **body["auto"]} if "schedule" in body and task.get("mode") == "scheduled": sch = {**task.get("schedule", {}), **body["schedule"]} try: sch["next_run"] = _schedule_next_run(sch) except ValueError as e: return jsonify({"error": f"调度配置错误: {e}"}), 400 task["schedule"] = sch task["updated_at"] = now_str() store.upsert_task(task) db.upsert_task(task) return jsonify(task) @app.route("/api/tasks/", methods=["DELETE"]) def api_delete_task(tid): """删除任务 -> 移入回收站 (可恢复)""" task = store.get_task(tid) if not task: return jsonify({"error": "任务不存在"}), 404 if task.get("deleted_at"): return jsonify({"error": "任务已在回收站中"}), 400 _stop_job(tid) with JOBS_LOCK: JOBS.pop(tid, None) store.soft_delete_task(tid, now_str()) db.upsert_task(store.get_task(tid)) return jsonify({"ok": True, "msg": "已移入回收站"}) # ---------------- API: 回收站 ---------------- def _purge_out_dir(out_dir): """删除任务输出目录; 仅当目录位于项目 out/ 下才删 (自定义目录保留), 返回是否删除""" try: d = os.path.realpath(out_dir) base = os.path.realpath(os.path.join(HERE, "out")) if d.startswith(base + os.sep) and os.path.isdir(d): import shutil shutil.rmtree(d, ignore_errors=True) return True except Exception: pass return False @app.route("/api/trash", methods=["GET"]) def api_trash_list(): items = store.list_trash() for t in items: t["runs_count"] = len(store.get_runs(t["id"])) t["out_dir"] = resolve_out_dir(t) return jsonify(items) @app.route("/api/trash//restore", methods=["POST"]) def api_trash_restore(tid): task = store.get_task(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)) return jsonify({"ok": True, "msg": "已恢复"}) @app.route("/api/trash/", methods=["DELETE"]) def api_trash_purge(tid): task = store.get_task(tid) if not task or not task.get("deleted_at"): return jsonify({"error": "任务不在回收站中"}), 404 out_dir = resolve_out_dir(task) removed = _purge_out_dir(out_dir) store.purge_task(tid) db.purge_task_db(tid) return jsonify({"ok": True, "purged": True, "files_removed": removed, "out_dir": out_dir}) @app.route("/api/trash", methods=["DELETE"]) def api_trash_clear(): items = store.list_trash() dirs = [resolve_out_dir(t) for t in items] store.purge_trash() db.purge_trash_db() removed = sum(1 for d in dirs if _purge_out_dir(d)) return jsonify({"ok": True, "purged": len(items), "dirs_removed": removed}) # ---------------- API: 运行控制 ---------------- @app.route("/api/tasks//start", methods=["POST"]) def api_start(tid): task = store.get_task(tid) if not task: return jsonify({"error": "任务不存在"}), 404 run, err = start_run(task) if err: return jsonify({"error": err}), 409 return jsonify(run) @app.route("/api/tasks//stop", methods=["POST"]) def api_stop(tid): if _stop_job(tid): return jsonify({"ok": True, "msg": "已发送终止信号"}) return jsonify({"ok": True, "msg": "任务未在运行"}) @app.route("/api/tasks//pause", methods=["POST"]) def api_pause(tid): with JOBS_LOCK: job = JOBS.get(tid) if job and job.is_running(): job.pause() return jsonify({"ok": True}) return jsonify({"error": "任务未在运行"}), 409 @app.route("/api/tasks//resume", methods=["POST"]) def api_resume(tid): with JOBS_LOCK: job = JOBS.get(tid) if job and job.is_running(): job.resume() return jsonify({"ok": True}) return jsonify({"error": "任务未在运行"}), 409 # ---------------- API: 统计 ---------------- @app.route("/api/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")) total_runs = ok = fail = imgs = 0 for t in tasks: for r in store.get_runs(t["id"]): total_runs += 1 st = r.get("stats") or {} ok += st.get("ok", 0) fail += st.get("fail", 0) 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 return jsonify({ "tasks": len(tasks), "running": running, "runs": total_runs, "ok": ok, "fail": fail, "images": imgs, "disk_mb": round(size / 1048576, 1), "trash": trash_count, }) # ---------------- API: 试爬取 ---------------- @app.route("/api/probe", methods=["POST"]) def api_probe(): """试爬取: 按表单给出的规则探测起始页, 返回将爬取的链接清单""" body = request.get_json(force=True) or {} seed = (body.get("seed_url") or "").strip() if not seed: return jsonify({"error": "请填写起始网址"}), 400 if not seed.startswith("http"): seed = "https://" + seed result = probe_links( seed, include=[x.strip() for x in (body.get("include") or []) if x.strip()], exclude=[x.strip() for x in (body.get("exclude") or []) if x.strip()], same_domain=body.get("same_domain", True), use_regex=bool(body.get("use_regex", False)), timeout=int(body.get("timeout") or 45), ) return jsonify(result) @app.route("/api/tasks//probe", methods=["POST"]) def api_task_probe(tid): """对已保存的自动任务执行试爬取 (使用保存的规则)""" task = store.get_task(tid) if not task: return jsonify({"error": "任务不存在"}), 404 if task.get("mode") != "auto": return jsonify({"error": "仅自动爬取任务支持试爬取"}), 400 auto = task.get("auto", {}) result = probe_links( auto.get("seed_url", ""), include=auto.get("include", []), exclude=auto.get("exclude", []), same_domain=auto.get("same_domain", True), use_regex=bool(auto.get("use_regex", False)), timeout=int((task.get("config") or {}).get("timeout", 60)), ) return jsonify(result) # ---------------- API: 搜索 ---------------- @app.route("/api/search") def api_search(): """基本搜索: 对任务名/网址/标题做子串匹配 (仅遍历本地内存数据, 轻量)""" q = (request.args.get("q") or "").strip() if not q: return jsonify({"q": q, "count": 0, "results": []}) ql = q.lower() results = [] seen = set() for task in store.load_tasks(): if task.get("deleted_at"): continue # 回收站任务不参与搜索 auto = task.get("auto") or {} task_hit = ( ql in (task.get("name") or "").lower() 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 r in run.get("results", []): url = r.get("url") or "" title = r.get("title") or "" if ql not in url.lower() and ql not in title.lower(): continue key = (task["id"], run.get("id"), url) if key in seen: continue seen.add(key) results.append({ "type": "page", "task_id": task["id"], "task_name": task.get("name", ""), "mode": task.get("mode", ""), "run_id": run.get("id", ""), "run_status": run.get("status", ""), "url": url, "title": title, "status": r.get("status", ""), "crawl_time": r.get("crawl_time", ""), "html_file": r.get("html_file", ""), "txt_file": r.get("txt_file", ""), "meta_file": r.get("meta_file", ""), "images": len(r.get("images", []) or []), }) if len(results) >= 100: break if len(results) >= 100: break if task_hit and len(results) < 100: results.append({ "type": "task", "task_id": task["id"], "task_name": task.get("name", ""), "mode": task.get("mode", ""), "urls_count": len(task.get("urls", [])), "created_at": task.get("created_at", ""), "seed_url": auto.get("seed_url", ""), }) return jsonify({"q": q, "count": len(results), "results": results}) # ---------------- API: 运行记录与文件 ---------------- @app.route("/api/runs/") def api_run_detail(rid): run = store.get_run(rid) if not run: return jsonify({"error": "运行记录不存在"}), 404 return jsonify(run) @app.route("/api/runs//logs") def api_run_logs(rid): run = store.get_run(rid) if not run: return jsonify({"error": "运行记录不存在"}), 404 offset = int(request.args.get("offset", 0)) logs = run.get("logs", []) return jsonify({"logs": logs[offset:], "count": len(logs)}) @app.route("/api/file") def api_file(): tid = request.args.get("task_id", "") path = request.args.get("path", "") task = store.get_task(tid) if not task or not path: return jsonify({"error": "参数错误"}), 400 out_dir = os.path.realpath(resolve_out_dir(task)) full = os.path.realpath(os.path.join(out_dir, path)) if not full.startswith(out_dir + os.sep) and full != out_dir: return jsonify({"error": "路径越界"}), 403 if not os.path.isfile(full): return jsonify({"error": "文件不存在"}), 404 return send_file(full) # ---------------- 启动 ---------------- scheduler = Scheduler(start_run) def main(): os.makedirs(os.path.join(HERE, "data"), exist_ok=True) os.makedirs(os.path.join(HERE, "out"), exist_ok=True) db.init_db() # 全量同步存量任务到数据库 for t in store.load_tasks(): db.upsert_task(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) if __name__ == "__main__": main()