# -*- 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