Files
universal-crawler/store.py
T

187 lines
4.6 KiB
Python

# -*- coding: utf-8 -*-
"""任务 / 运行记录持久化层 (JSON 文件存储, 线程安全)"""
import json
import os
import threading
import uuid
from datetime import datetime
HERE = os.path.dirname(os.path.abspath(__file__))
DATA_DIR = os.path.join(HERE, "data")
TASKS_FILE = os.path.join(DATA_DIR, "tasks.json")
RUNS_FILE = os.path.join(DATA_DIR, "runs.json") # {task_id: [run, ...]}
_lock = threading.RLock()
def new_id(prefix):
return f"{prefix}_{uuid.uuid4().hex[:12]}"
def now_str():
return datetime.now().strftime("%Y-%m-%d %H:%M:%S")
def _load(path, default):
if not os.path.exists(path):
return default
try:
with open(path, encoding="utf-8") as f:
return json.load(f)
except Exception:
return default
def _save(path, data):
os.makedirs(os.path.dirname(path), exist_ok=True)
tmp = path + ".tmp"
with open(tmp, "w", encoding="utf-8") as f:
json.dump(data, f, ensure_ascii=False, indent=2)
os.replace(tmp, path)
# ---------------- 任务 ----------------
def load_tasks():
with _lock:
return _load(TASKS_FILE, [])
def get_task(task_id):
for t in load_tasks():
if t["id"] == task_id:
return t
return None
def upsert_task(task):
with _lock:
tasks = _load(TASKS_FILE, [])
for i, t in enumerate(tasks):
if t["id"] == task["id"]:
tasks[i] = task
break
else:
tasks.append(task)
_save(TASKS_FILE, tasks)
def delete_task(task_id):
"""彻底删除任务 (任务 + 运行记录)"""
with _lock:
tasks = _load(TASKS_FILE, [])
tasks = [t for t in tasks if t["id"] != task_id]
_save(TASKS_FILE, tasks)
runs = _load(RUNS_FILE, {})
runs.pop(task_id, None)
_save(RUNS_FILE, runs)
# ---------------- 回收站 ----------------
def soft_delete_task(task_id, ts):
"""删除任务 -> 移入回收站 (软删除, 记录保留可恢复)"""
with _lock:
tasks = _load(TASKS_FILE, [])
for t in tasks:
if t["id"] == task_id:
t["deleted_at"] = ts
break
_save(TASKS_FILE, tasks)
def restore_task(task_id):
"""从回收站恢复任务"""
with _lock:
tasks = _load(TASKS_FILE, [])
for t in tasks:
if t["id"] == task_id:
t.pop("deleted_at", None)
break
_save(TASKS_FILE, tasks)
def list_trash():
"""回收站任务列表 (按删除时间倒序)"""
with _lock:
tasks = _load(TASKS_FILE, [])
return sorted(
[t for t in tasks if t.get("deleted_at")],
key=lambda x: x.get("deleted_at", ""), reverse=True,
)
def purge_task(task_id):
"""从回收站彻底删除单个任务 (任务 + 运行记录)"""
with _lock:
tasks = _load(TASKS_FILE, [])
tasks = [t for t in tasks if t["id"] != task_id]
_save(TASKS_FILE, tasks)
runs = _load(RUNS_FILE, {})
runs.pop(task_id, None)
_save(RUNS_FILE, runs)
def purge_trash():
"""清空回收站, 返回被清空的任务 id 列表"""
with _lock:
tasks = _load(TASKS_FILE, [])
kept, purged = [], []
for t in tasks:
if t.get("deleted_at"):
purged.append(t["id"])
else:
kept.append(t)
_save(TASKS_FILE, kept)
runs = _load(RUNS_FILE, {})
for tid in purged:
runs.pop(tid, None)
_save(RUNS_FILE, runs)
return purged
# ---------------- 运行记录 ----------------
def load_runs_map():
with _lock:
return _load(RUNS_FILE, {})
def get_runs(task_id):
with _lock:
return _load(RUNS_FILE, {}).get(task_id, [])
def add_run(task_id, run):
with _lock:
runs = _load(RUNS_FILE, {})
runs.setdefault(task_id, []).append(run)
if len(runs[task_id]) > 50: # 每个任务最多保留 50 次运行记录
del runs[task_id][:-50]
_save(RUNS_FILE, runs)
def save_run(task_id, run):
with _lock:
runs = _load(RUNS_FILE, {})
lst = runs.setdefault(task_id, [])
for i, r in enumerate(lst):
if r["id"] == run["id"]:
lst[i] = run
break
else:
lst.append(run)
if len(lst) > 50:
del lst[:-50]
_save(RUNS_FILE, runs)
def get_run(run_id):
with _lock:
runs = _load(RUNS_FILE, {})
for lst in runs.values():
for r in lst:
if r["id"] == run_id:
return r
return None