# coding: utf-8
"""
process.py —— torrent 下载中央监控 / 控制服务（固定机器上运行）
================================================================

部署：在固定机器（IP 21.196.89.120）上执行  `python process.py`  起一个 HTTP 服务，
汇总所有 downloader 上报的状态并展示，同时可通过接口控制 downloader 任务的
暂停 / 开始 / 删除。

持久化
------
所有任务快照写入本机 SQLite（Approach B：downloader 只负责 HTTP 上报，process 独占
DB 单写入端）。进程重启时从 DB 加载全部记录，看板不因重启丢历史。DB 路径由
环境变量 TORRENT_PROCESS_DB 指定（默认与本文件同目录 process_state.db）。

HTTP 接口
---------
  POST /api/report      downloader 周期性上报本机某个种子的状态快照（JSON）。
                        响应体返回该任务当前“期望状态” {"desired": "running|paused|deleted"}。
  POST /api/control     控制某个任务：body {"task_key": "rid:info_hash", "action": "pause|start|delete|force_ingest"}
  POST /api/retry       重试：body {"task_key": "..."} 或 {"task_keys": [...]}，把原始种子入队到原机本机重试队列（复用本地文件重跑，不走远端 job_id）
  GET  /api/pending_retries  原机 downloader 轮询取回本机待重试的原始种子：?machine=X -> {"seeds": [...]}
  GET  /api/summary     过滤后的分类计数 / 机器磁盘 / 速度统计 / requirement_id 候选。
  GET  /api/list        某分类的一页轻量行（不含 files/bt，支持 offset/limit 无限下拉）。
  GET  /api/task        单个 task_key 的完整明细（files/bt，展开时按需拉取）。
  GET  /api/export      按当前筛选（requirement_id + 时间）导出某分类为 CSV。
  GET  /api/status      全量状态（兼容旧接口 / 控制台）。
  GET  /api/dispatch    重复投递决策（completed/ignore/proceed + mode）。
  GET  /                Web 看板（分类无限下拉懒加载 + 按 requirement_id/添加时间筛选 +
                        分类导出 + 静默刷新防卡顿）。

筛选：requirement_id 为必填字段（downloader 必须上报，看板按其 + 添加时间联合筛选）。
"""
from __future__ import annotations

import io
import json
import os
import sqlite3
import threading
import time
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
from typing import Any, Dict, List, Optional, Tuple
from urllib.parse import urlparse, parse_qs

import requests
from loguru import logger


# ============================================================================
# 常量
# ============================================================================

# 固定监控机器
PROCESS_SERVER_HOST = "21.196.89.120"
PROCESS_SERVER_PORT = 8686
PROCESS_SERVER_URL = f"http://{PROCESS_SERVER_HOST}:{PROCESS_SERVER_PORT}"

# 持久化 DB 路径（process 独占单写）
DB_PATH = os.environ.get(
    "TORRENT_PROCESS_DB",
    os.path.join(os.path.dirname(os.path.abspath(__file__)), "process_state.db"))
# 后台刷盘间隔：progress-only 变更攒批写；状态/期望态/新任务立即写
FLUSH_INTERVAL = float(os.environ.get("TORRENT_PROCESS_FLUSH_INTERVAL", "5"))

# 文件 / 种子状态
ST_PENDING = "pending"
ST_METADATA = "metadata"          # 正在获取种子元数据（尚未开始下载文件）
ST_CHECKING = "checking"          # 正在校验磁盘上已存在的数据（本地哈希校验，非下载）
ST_DOWNLOADING = "downloading"
ST_UPLOADING = "uploading"
ST_COMPLETED = "已完成"
ST_ERROR = "error"
ST_SKIPPED = "skipped"
ST_PAUSED = "paused"
ST_DELETED = "deleted"
ST_STALLED = "stalled"                 # 下载停滞/未完成被放弃（仍有文件没下完）
ST_EXTRACT_FAILED = "extract_failed"   # 已全部下完但解压/入库失败
ST_FORCE_INGESTED = "force_ingested"   # 已强制入库（只处理了已下载完成的文件）
ST_SEEDING = "seeding"                 # 下载完成后持续做种/分享（模仿 qB）
ST_CANCELED_NO_SPACE = "canceled_no_space"   # 种子过大、整机都塞不下，取消下载
ST_DELETED_NO_SPACE = "deleted_no_space"     # 已下载入库后因存储不足被 FIFO 删除
ST_REQUEUED = "requeued"   # 已点“重试”：原始种子已入队到原机本机重试队列，等原机本地重跑
ST_INGESTED_DONE = "ingested_done"   # 下载+做种结束、进程已释放：数据已全部入库(COS)，本地已清
ST_QUEUED = "queued"   # 已投递、在本机线程池排队等待（>max_worker 的种子尚未开始处理）

# 下载器自身 job_id：无进度(无本地数据)的重试可直接重投到这里，由调度分配到任意机器
# 重新下载（有进度的则按原机路由复用本地文件，见 MonitorStore.retry）。
RETRY_JOB_ID = os.environ.get("TORRENT_RETRY_JOB_ID",
                              "16441f134ef940cc89305d78e4e9bb04")

# 状态 -> 中文展示文案（看板用）
STATE_LABELS = {
    ST_PENDING: "等待中",
    ST_METADATA: "获取元数据中",
    ST_CHECKING: "校验中",
    ST_DOWNLOADING: "下载中",
    ST_UPLOADING: "上传中",
    ST_COMPLETED: "已完成",
    ST_ERROR: "错误",
    ST_SKIPPED: "已跳过",
    ST_PAUSED: "已暂停",
    ST_DELETED: "已删除",
    ST_STALLED: "下载停滞放弃",
    ST_EXTRACT_FAILED: "解压/入库失败",
    ST_FORCE_INGESTED: "已强制入库",
    ST_SEEDING: "做种中",
    ST_CANCELED_NO_SPACE: "因存储不足取消下载",
    ST_DELETED_NO_SPACE: "因存储不足已下载入库后删除",
    ST_REQUEUED: "重投队列",
    ST_INGESTED_DONE: "已全部入库完成",
    ST_QUEUED: "排队等待下载",
}

# 期望状态（控制指令的落地态）
DESIRED_RUNNING = "running"
DESIRED_PAUSED = "paused"
DESIRED_DELETED = "deleted"
# 强制入库：不再等待整个种子下载完成，直接把“已下载完成的文件”处理并入库后结束
DESIRED_FORCE_INGEST = "force_ingest"

# 任务快照超过该秒数无上报视为离线
OFFLINE_SECONDS = 60

# 看板分类顺序（key）。做种中放最后（页面最底部）；重投队列在其上。
GROUP_KEYS = ["queued", "downloading", "paused", "stalled", "extract_failed",
              "force_ingested", "canceled_no_space", "deleted_no_space", "deleted",
              "other", "requeued", "seeding", "ingested_done"]


def _fmt_size(n: float) -> str:
    n = float(n or 0)
    for unit in ("B", "KB", "MB", "GB", "TB"):
        if n < 1024 or unit == "TB":
            return f"{n:.1f}{unit}"
        n /= 1024
    return f"{n:.1f}TB"


def _fmt_speed(bps: float) -> str:
    return _fmt_size(bps) + "/s"


def _category_of(state: str, desired: str, offline: bool = False) -> str:
    """把 (state, desired) 归到看板分类（与前端保持一致，供 list/summary/export 共用）。

    未识别的状态一律落入 "other"（其他特殊情况），保证任何种子都能在前端显示、不丢失。
    离线判定：下载中/元数据/校验/上传/等待 等“进行中”状态若掉线（长时间无上报），
    归入 "other"（其他特殊情况）——机器/进程异常导致的中断，需人工关注/重试。
    """
    if state == ST_QUEUED:
        return "queued"
    if state == ST_SEEDING:
        # 做种中但已离线（做种进程已退出/失联）：数据已入库 COS，视为“已全部入库完成”
        return "ingested_done" if offline else "seeding"
    if state == ST_INGESTED_DONE:
        return "ingested_done"
    if state == ST_REQUEUED:
        return "requeued"
    if state == ST_FORCE_INGESTED:
        return "force_ingested"
    if state == ST_CANCELED_NO_SPACE:
        return "canceled_no_space"
    if state == ST_DELETED_NO_SPACE:
        return "deleted_no_space"
    if state == ST_EXTRACT_FAILED:
        return "extract_failed"
    if state in (ST_STALLED, ST_ERROR):
        return "stalled"
    if state == ST_DELETED or desired == DESIRED_DELETED:
        return "deleted"
    if desired == DESIRED_PAUSED or state == ST_PAUSED:
        return "paused"
    if state in (ST_DOWNLOADING, ST_METADATA, ST_CHECKING, ST_UPLOADING,
                 ST_PENDING, ST_COMPLETED, ST_SKIPPED):
        # 进行中/已完成待交接的任务若掉线（长时间无上报），一律归入 "other"
        # （其他特殊情况）——机器/进程异常导致的中断，需人工关注/重试；否则会顶着
        # “离线”徽标仍堆在“正在下载”分组里。
        return "other" if offline else "downloading"
    return "other"


def _task_key(rid, info_hash: str) -> str:
    """任务主键：同 info_hash 不同 requirement_id 属不同任务，避免互相覆盖/漏显示。"""
    try:
        r = int(rid or 0)
    except (TypeError, ValueError):
        r = 0
    return f"{r}:{info_hash}"


def _to_float(v) -> Optional[float]:
    try:
        if v is None or v == "":
            return None
        return float(v)
    except (TypeError, ValueError):
        return None


# ============================================================================
# 服务端状态存储（内存为读源 + SQLite 持久化）
# ============================================================================

class MonitorStore:
    """线程安全的任务状态 + 控制期望态存储，写穿到 SQLite。"""

    def __init__(self, db_path: str = DB_PATH):
        self._lock = threading.RLock()
        # task_key -> snapshot(dict)（downloader 上报的最新快照，含 machine/files 等）
        self._status: Dict[str, dict] = {}
        # task_key -> desired 状态
        self._desired: Dict[str, str] = {}
        # info_hash -> 首次出现（添加）时间戳（用于按添加时间筛选，不被覆盖）
        self._added_ts: Dict[str, float] = {}
        # 待刷盘（progress-only 变更攒批）
        self._dirty: set = set()
        # 本机重试队列：machine -> [raw_seed,...]。重试改为“原机本地重投”（复用本地已下载
        # 文件、尝试解压入库），不走远端 job_id（避免被调度到别的机器重新下载一遍）。
        self._retry_queue: Dict[str, list] = {}

        self.db_path = db_path
        self._db = sqlite3.connect(db_path, check_same_thread=False)
        self._db.execute("PRAGMA journal_mode=WAL")
        self._db.execute("PRAGMA synchronous=NORMAL")
        self._init_schema()
        self._load_all()
        self._load_retry_queue()

        self._stop = threading.Event()
        threading.Thread(target=self._flush_loop, daemon=True,
                         name="store-flush").start()

    # -------------------------------------------------- 持久化底层
    def _migrate_legacy_tasks(self):
        """旧版以 info_hash 为主键的 tasks 表迁移到 task_key 复合主键。"""
        logger.info("[process] 迁移 tasks 表：info_hash 主键 -> task_key 复合主键")
        self._db.execute("ALTER TABLE tasks RENAME TO tasks_legacy")
        self._db.execute("""
            CREATE TABLE tasks (
                task_key TEXT PRIMARY KEY,
                info_hash TEXT,
                requirement_id INTEGER,
                name TEXT,
                machine TEXT,
                state TEXT,
                desired TEXT,
                category TEXT,
                progress REAL,
                download_rate REAL,
                upload_rate REAL,
                num_seeds INTEGER,
                num_peers INTEGER,
                num_files INTEGER,
                completed_files INTEGER,
                total_size INTEGER,
                added_ts REAL,
                updated_ts REAL,
                snapshot_json TEXT
            )""")
        legacy_cols = {r[1] for r in self._db.execute(
            "PRAGMA table_info(tasks_legacy)")}
        if "snapshot_json" in legacy_cols:
            self._db.execute("""
                INSERT INTO tasks(task_key,info_hash,requirement_id,name,machine,state,
                    desired,category,progress,download_rate,upload_rate,num_seeds,
                    num_peers,num_files,completed_files,total_size,added_ts,updated_ts,
                    snapshot_json)
                SELECT printf('%s:%s', COALESCE(requirement_id, 0), info_hash),
                    info_hash, requirement_id, name, machine, state, desired, category,
                    progress, download_rate, upload_rate, num_seeds, num_peers, num_files,
                    completed_files, total_size, added_ts, updated_ts, snapshot_json
                FROM tasks_legacy""")
        self._db.execute("DROP TABLE tasks_legacy")
        self._db.commit()

    def _init_schema(self):
        cur = self._db.execute(
            "SELECT name FROM sqlite_master WHERE type='table' AND name='tasks'")
        if cur.fetchone():
            cols = {r[1] for r in self._db.execute("PRAGMA table_info(tasks)")}
            if "task_key" not in cols:
                self._migrate_legacy_tasks()
        self._db.execute("""
            CREATE TABLE IF NOT EXISTS tasks (
                task_key TEXT PRIMARY KEY,
                info_hash TEXT,
                requirement_id INTEGER,
                name TEXT,
                machine TEXT,
                state TEXT,
                desired TEXT,
                category TEXT,
                progress REAL,
                download_rate REAL,
                upload_rate REAL,
                num_seeds INTEGER,
                num_peers INTEGER,
                num_files INTEGER,
                completed_files INTEGER,
                total_size INTEGER,
                added_ts REAL,
                updated_ts REAL,
                snapshot_json TEXT
            )""")
        self._db.execute(
            "CREATE INDEX IF NOT EXISTS idx_rid_added ON tasks(requirement_id, added_ts)")
        self._db.execute("CREATE INDEX IF NOT EXISTS idx_state ON tasks(state)")
        # 错误/告警事件表（解压失败、存储不足跳过、入库失败、停滞放弃等），持久化供
        # 前端“消息列表”查询历史。
        self._db.execute("""
            CREATE TABLE IF NOT EXISTS events (
                id INTEGER PRIMARY KEY AUTOINCREMENT,
                ts REAL,
                info_hash TEXT,
                requirement_id INTEGER,
                machine TEXT,
                name TEXT,
                level TEXT,
                message TEXT
            )""")
        self._db.execute(
            "CREATE INDEX IF NOT EXISTS idx_ev_ts ON events(ts)")
        self._db.execute(
            "CREATE INDEX IF NOT EXISTS idx_ev_rid_ts ON events(requirement_id, ts)")
        # 本机重试队列（持久化：process 重启后原机仍能取回本机待重试种子）。
        self._db.execute("""
            CREATE TABLE IF NOT EXISTS retry_queue (
                id INTEGER PRIMARY KEY AUTOINCREMENT,
                machine TEXT,
                raw_seed TEXT,
                ts REAL
            )""")
        self._db.execute(
            "CREATE INDEX IF NOT EXISTS idx_retry_machine ON retry_queue(machine)")
        self._db.commit()

    def _load_all(self):
        cur = self._db.execute(
            "SELECT task_key, info_hash, requirement_id, desired, added_ts, "
            "updated_ts, snapshot_json FROM tasks")
        n = 0
        for key, ih, rid, desired, added_ts, updated_ts, snap_json in cur.fetchall():
            try:
                snap = json.loads(snap_json) if snap_json else {}
            except Exception:
                snap = {}
            snap["info_hash"] = ih
            snap["requirement_id"] = int(rid or 0)
            snap["_added_ts"] = float(added_ts or 0)
            # 重启后按最后一次更新时间计算离线（未再上报即显示离线）
            snap["_report_ts"] = float(updated_ts or 0)
            self._status[key] = snap
            self._desired[key] = desired or DESIRED_RUNNING
            self._added_ts[key] = float(added_ts or 0)
            n += 1
        logger.info(f"[process] 已从 DB 加载 {n} 条历史任务记录 ({self.db_path})")

    def _load_retry_queue(self):
        """加载持久化的本机重试队列（process 重启后原机仍能取回）。"""
        n = 0
        try:
            cur = self._db.execute(
                "SELECT machine, raw_seed FROM retry_queue ORDER BY id")
            for machine, raw in cur.fetchall():
                if machine and raw:
                    self._retry_queue.setdefault(machine, []).append(raw)
                    n += 1
        except Exception as e:
            logger.warning(f"[process] 加载重试队列失败: {e}")
        if n:
            logger.info(f"[process] 已从 DB 加载 {n} 条待重试种子")

    def _persist_locked(self, key: str):
        snap = self._status.get(key)
        if snap is None:
            return
        desired = self._desired.get(key, DESIRED_RUNNING)
        state = snap.get("state", "")
        cat = _category_of(state, desired)
        try:
            self._db.execute(
                "INSERT OR REPLACE INTO tasks(task_key,info_hash,requirement_id,name,"
                "machine,state,desired,category,progress,download_rate,upload_rate,"
                "num_seeds,num_peers,num_files,completed_files,total_size,added_ts,"
                "updated_ts,snapshot_json) VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)",
                (key, str(snap.get("info_hash") or ""),
                 int(snap.get("requirement_id") or 0),
                 str(snap.get("name") or ""), str(snap.get("machine") or ""),
                 state, desired, cat, float(snap.get("progress") or 0),
                 float(snap.get("download_rate") or 0),
                 float(snap.get("upload_rate") or 0),
                 int(snap.get("num_seeds") or 0), int(snap.get("num_peers") or 0),
                 int(snap.get("num_files") or 0), int(snap.get("completed_files") or 0),
                 int(snap.get("total_size") or 0),
                 float(self._added_ts.get(key, 0)),
                 float(snap.get("_report_ts", 0)),
                 json.dumps(snap, ensure_ascii=False)))
            self._db.commit()
            self._dirty.discard(key)
        except Exception as e:
            logger.warning(f"[process] 持久化失败 key={key}: {e}")

    def _delete_locked(self, key: str):
        try:
            self._db.execute("DELETE FROM tasks WHERE task_key=?", (key,))
            self._db.commit()
        except Exception as e:
            logger.warning(f"[process] 删除持久化记录失败 key={key}: {e}")

    def _insert_events_locked(self, events: list, info_hash: str, rid: int,
                              machine: str, name: str):
        """把 downloader 本轮上报的错误/告警事件写入 events 表（持久化）。"""
        try:
            rows = []
            for ev in events:
                if not isinstance(ev, dict):
                    continue
                rows.append((
                    float(ev.get("ts") or time.time()), info_hash, int(rid or 0),
                    machine, name, str(ev.get("level") or "info"),
                    str(ev.get("message") or "")[:1000]))
            if rows:
                self._db.executemany(
                    "INSERT INTO events(ts,info_hash,requirement_id,machine,name,"
                    "level,message) VALUES(?,?,?,?,?,?,?)", rows)
                self._db.commit()
        except Exception as e:
            logger.warning(f"[process] 事件落库失败 info_hash={info_hash}: {e}")

    def list_events(self, rid=None, start: Optional[float] = None,
                    end: Optional[float] = None, level: str = "",
                    offset: int = 0, limit: int = 30, machine=None, q=None) -> dict:
        """按 requirement_id + 时间范围（事件发生时间）+ 级别 + 机器 + 搜索 分页查询历史事件。

        搜索 q：匹配 info_hash / 种子名称(name) / 错误信息(message)，用于“按种子或按
        错误内容检索报错消息”。
        """
        where = []
        args: list = []
        if rid not in (None, "", "all"):
            try:
                where.append("requirement_id=?"); args.append(int(rid))
            except (TypeError, ValueError):
                pass
        if start is not None:
            where.append("ts>=?"); args.append(float(start))
        if end is not None:
            where.append("ts<=?"); args.append(float(end))
        if level:
            where.append("level=?"); args.append(str(level))
        if machine not in (None, "", "all"):
            where.append("machine=?"); args.append(str(machine))
        if q:
            ql = str(q).strip()
            if ql:
                like = f"%{ql}%"
                where.append(
                    "(info_hash LIKE ? OR name LIKE ? OR message LIKE ? OR info_hash=?)")
                args.extend([like, like, like, ql])
        wsql = (" WHERE " + " AND ".join(where)) if where else ""
        rows = []
        total = 0
        with self._lock:
            try:
                total = self._db.execute(
                    "SELECT COUNT(*) FROM events" + wsql, args).fetchone()[0]
                cur = self._db.execute(
                    "SELECT ts,info_hash,requirement_id,machine,name,level,message "
                    "FROM events" + wsql + " ORDER BY ts DESC, id DESC LIMIT ? OFFSET ?",
                    args + [int(limit) if limit else -1, int(offset)])
                for ts, ih, r, mac, nm, lv, msg in cur.fetchall():
                    rows.append({
                        "ts": ts, "info_hash": ih, "requirement_id": r,
                        "machine": mac, "name": nm, "level": lv, "message": msg,
                    })
            except Exception as e:
                logger.warning(f"[process] 查询事件失败: {e}")
        return {"rows": rows, "total": total, "offset": offset, "limit": limit}

    def _flush_loop(self):
        while not self._stop.wait(FLUSH_INTERVAL):
            try:
                with self._lock:
                    dirty = list(self._dirty)
                for ih in dirty:
                    with self._lock:
                        self._persist_locked(ih)
            except Exception as e:
                logger.warning(f"[process] 刷盘循环异常: {e}")

    # -------------------------------------------------- 上报 / 控制
    def report(self, snapshot: dict) -> str:
        info_hash = str(snapshot.get("info_hash") or "")
        if not info_hash:
            return DESIRED_RUNNING
        try:
            rid = int(snapshot.get("requirement_id") or 0)
        except (TypeError, ValueError):
            rid = 0
        if not rid:
            logger.warning(f"[process] 上报缺少 requirement_id（必填）info_hash={info_hash} "
                           f"machine={snapshot.get('machine')}")
        snapshot["requirement_id"] = rid
        key = _task_key(rid, info_hash)
        # 取出本轮事件（不随 snapshot_json 落库，单独进 events 表持久化）
        events = snapshot.pop("events", None) or []
        now = time.time()
        with self._lock:
            if events:
                self._insert_events_locked(
                    events, info_hash, rid,
                    str(snapshot.get("machine") or ""),
                    str(snapshot.get("name") or ""))
            prev = self._status.get(key)
            # 做种/停做种等阶段的快照不带文件明细（避免每 20s 重发上万条文件）：沿用上一份
            # 快照的 files / num_files / completed_files / bt，避免详情文件列表从满变空、
            # “文件 0/0”。仅在新快照确实缺这些时才补，正常下载快照自带则不影响。
            if prev is not None:
                if not snapshot.get("files") and prev.get("files"):
                    snapshot["files"] = prev["files"]
                if not snapshot.get("num_files") and prev.get("num_files"):
                    snapshot["num_files"] = prev["num_files"]
                if not snapshot.get("completed_files") and prev.get("completed_files"):
                    snapshot["completed_files"] = prev["completed_files"]
                if not snapshot.get("bt") and prev.get("bt"):
                    snapshot["bt"] = prev["bt"]
            if key not in self._added_ts:
                self._added_ts[key] = now
            snapshot["_report_ts"] = now
            snapshot["_added_ts"] = self._added_ts[key]
            self._status[key] = snapshot
            desired = self._desired.get(key, DESIRED_RUNNING)
            # 状态变化 / 新任务立即落盘；纯进度变更攒批（后台刷盘）
            state_changed = (prev is None) or \
                (prev.get("state") != snapshot.get("state"))
            self._dirty.add(key)
            if state_changed:
                self._persist_locked(key)
            return desired

    def set_control(self, task_key: str, action: str) -> bool:
        action = (action or "").lower()
        mapping = {
            "pause": DESIRED_PAUSED,
            "start": DESIRED_RUNNING,
            "resume": DESIRED_RUNNING,
            "delete": DESIRED_DELETED,
            "force_ingest": DESIRED_FORCE_INGEST,
            "force": DESIRED_FORCE_INGEST,
        }
        desired = mapping.get(action)
        if desired is None:
            return False
        with self._lock:
            self._desired[task_key] = desired
            snap = self._status.get(task_key)
            # 强制入库：若任务已停滞放弃/错误/离线（没有在跑的下载实例来消费该指令），
            # 直接把状态置为“已强制入库”移入对应分组。停滞任务里已下载完成的文件在下载
            # 过程中已逐个投种入库，剩下的只是失败/未完成文件，force_ingest 语义即“接受
            # 已入库的部分、结束该任务”，故直接落终态即可，无需等待不存在的下载实例。
            if desired == DESIRED_FORCE_INGEST and snap is not None:
                now = time.time()
                state = snap.get("state", "")
                offline = (now - snap.get("_report_ts", 0)) > OFFLINE_SECONDS
                if offline or state in (ST_STALLED, ST_ERROR):
                    snap["state"] = ST_FORCE_INGESTED
                    snap["error"] = ""
                    snap["_report_ts"] = now
            if task_key in self._status:
                self._persist_locked(task_key)
        logger.info(f"[process] 控制指令 key={task_key} action={action} "
                    f"-> desired={desired}")
        return True

    def drop_task(self, task_key: str) -> bool:
        """从看板与 DB 移除一个任务记录（用于排队占位在真正开始下载后清除）。"""
        if not task_key:
            return False
        with self._lock:
            existed = task_key in self._status
            self._status.pop(task_key, None)
            self._desired.pop(task_key, None)
            self._added_ts.pop(task_key, None)
            self._dirty.discard(task_key)
            self._delete_locked(task_key)
        return existed

    def _route_to_origin_machine_locked(self, key: str, machine: str,
                                        raw: str, now: float) -> None:
        """把原始种子入队到“原下载机器”的本机重试队列（须在持锁下调用）。

        用于：上游把种子重投到了一台没处理过它的机器 B，但原机 A 本地已有未下完的部分
        文件——交回 A 续传（复用本地文件），B 不重复下载。原机若失联，滞留超时后由
        reroute_stale_retries 兜底改投远端任意机器。
        """
        q = self._retry_queue.setdefault(machine, [])
        if raw not in q:                     # 去重，避免重复投递把同一种子堆多份
            q.append(raw)
            try:
                self._db.execute(
                    "INSERT INTO retry_queue(machine, raw_seed, ts) VALUES(?,?,?)",
                    (machine, raw, now))
                self._db.commit()
            except Exception as e:
                logger.warning(f"[process] 交回原机续传入队失败: {e}")
        snap = self._status.get(key)
        if snap is not None:
            snap["state"] = ST_REQUEUED
            snap["error"] = ""
            snap["_report_ts"] = now
            self._persist_locked(key)

    def dispatch_decision(self, info_hash: str, rid=0, caller_machine="") -> dict:
        """重复投递同一 (requirement_id, info_hash) 时的处理决策（下载端开始前询问）。

        去重按 (rid, info_hash)：同一需求下重复投递才 ignore/completed；不同需求复用
        同一种子各自独立处理（本机已有完整数据则 downloader 走本地直接入库快路径）。
        规则：
          - 本 rid 已完成                          -> completed（跳过）
          - 本 rid 在线下载中/元数据/上传/等待      -> ignore（存活实例在跑）
          - 本 rid 在线已暂停（存活实例）           -> ignore + desired=running（让其继续）
          - 原机 A 有未下完的本地文件(有进度)，且请求来自别的机器 B
                                                   -> 交回 A 续传队列，B ignore（复用 A 本地文件）
          - 停滞放弃 / 错误（无进度）              -> proceed, mode=restart
          - 掉线的下载中/元数据/上传/等待（无进度）-> proceed, mode=continue
          - 已删除 / 掉线-已完成残留 / 未见过      -> proceed
        """
        if not info_hash:
            return {"decision": "proceed", "reason": "no-hash", "mode": "new"}
        key = _task_key(rid, info_hash)
        now = time.time()
        # 有未下完本地数据、可交回原机续传的状态集合
        _resumable = {ST_STALLED, ST_ERROR, ST_EXTRACT_FAILED, ST_DOWNLOADING,
                      ST_METADATA, ST_CHECKING, ST_UPLOADING, ST_PENDING}
        with self._lock:
            snap = self._status.get(key)
            mode = "new"
            if snap is not None:
                state = snap.get("state", "")
                offline = (now - snap.get("_report_ts", 0)) > OFFLINE_SECONDS
                origin = snap.get("machine") or ""
                raw = snap.get("raw_seed") or ""
                has_progress = (float(snap.get("progress") or 0) > 0
                                or int(snap.get("completed_files") or 0) > 0)
                if state == ST_COMPLETED and not offline:
                    return {"decision": "completed", "reason": "completed",
                            "mode": "skip"}
                if not offline and state in (ST_DOWNLOADING, ST_METADATA,
                                             ST_CHECKING, ST_UPLOADING, ST_PENDING):
                    return {"decision": "ignore", "reason": f"downloading:{state}",
                            "mode": "continue"}
                if not offline and state == ST_PAUSED:
                    self._desired[key] = DESIRED_RUNNING
                    self._persist_locked(key)
                    return {"decision": "ignore", "reason": "resume-paused",
                            "mode": "continue"}
                # 原机 A 有未下完的本地文件(有进度)，且这次请求来自“别的机器 B”
                # （caller_machine 已知且 != 原机）：交回 A 的本机续传队列，B 不处理。
                # caller 未知或就是原机自己时不触发（避免 A 自恢复时自我路由死循环）。
                if (has_progress and origin and raw and caller_machine
                        and origin != caller_machine and state in _resumable):
                    self._route_to_origin_machine_locked(key, origin, raw, now)
                    return {"decision": "ignore", "reason": "resume-on-origin",
                            "mode": "continue"}
                # 到这里都会 proceed，仅区分 continue / restart
                if state in (ST_STALLED, ST_ERROR):
                    mode = "restart"
                elif state in (ST_DOWNLOADING, ST_METADATA, ST_CHECKING,
                               ST_UPLOADING, ST_PENDING) and offline:
                    mode = "continue"   # 曾在下载但掉线：续传本地已有部分
                else:
                    mode = "restart"
            # proceed：占位 + 复位 desired
            if key not in self._added_ts:
                self._added_ts[key] = now
            self._desired[key] = DESIRED_RUNNING
            self._status[key] = {
                "info_hash": info_hash, "state": ST_PENDING, "name": "",
                "machine": "", "requirement_id": int(rid or 0),
                "_report_ts": now, "_added_ts": self._added_ts[key],
                "files": [],
            }
            self._persist_locked(key)
            reason = "restart" if snap is not None else "new"
            return {"decision": "proceed", "reason": reason, "mode": mode}

    def retry(self, task_keys: List[str]) -> dict:
        """重试，按“是否有下载进度”决定路由：

        - **有进度**（progress>0 或已完成部分文件）：原机本地有部分下载文件，入队到**原
          下载机器**的本机重试队列，让原机复用本地文件续传/解压入库，避免同一种子在两台
          机器各下一份。
        - **无进度**（没下过、无本地数据）：直接重投到下载器 job_id(RETRY_JOB_ID)，由调度
          分配到任意机器重新下载（不必绑原机）。

        原始种子(raw_seed) 来自上报快照（被 write_seed_json_file 包装前的 seed 原文）。
        """
        now = time.time()
        queued_local = 0
        missing = 0
        remote_items: List[Tuple[str, str]] = []   # (task_key, raw_seed)
        with self._lock:
            for tk in task_keys:
                snap = self._status.get(tk)
                if not snap:
                    continue
                raw = snap.get("raw_seed") or ""
                if not raw:
                    missing += 1
                    continue
                machine = snap.get("machine") or ""
                has_progress = (float(snap.get("progress") or 0) > 0
                                or int(snap.get("completed_files") or 0) > 0)
                if has_progress and machine:
                    # 有本地下载数据 -> 原机本地重跑（复用文件）
                    self._retry_queue.setdefault(machine, []).append(str(raw))
                    try:
                        self._db.execute(
                            "INSERT INTO retry_queue(machine, raw_seed, ts) VALUES(?,?,?)",
                            (machine, str(raw), now))
                    except Exception as e:
                        logger.warning(f"[process] 重试队列持久化失败: {e}")
                    snap["state"] = ST_REQUEUED
                    snap["error"] = ""
                    snap["_report_ts"] = now
                    self._persist_locked(tk)
                    queued_local += 1
                else:
                    # 无进度（或无机器记录）-> 重投远端 job_id，任意机器重新下载
                    remote_items.append((tk, str(raw)))
            if queued_local:
                try:
                    self._db.commit()
                except Exception:
                    pass

        # 无进度的：批量重投到下载器 job_id
        remote_n = 0
        if remote_items:
            try:
                try:
                    from .common import XunLiuClient, split_by_size
                except ImportError:
                    from common import XunLiuClient, split_by_size
                client = XunLiuClient()
                seeds = [raw for _tk, raw in remote_items]
                for chunk in split_by_size(seeds, 40 * 1024 * 1024):
                    client.upload_seed(RETRY_JOB_ID, chunk)
                remote_n = len(seeds)
                with self._lock:
                    for tk, _raw in remote_items:
                        snap = self._status.get(tk)
                        if snap:
                            snap["state"] = ST_REQUEUED
                            snap["error"] = ""
                            snap["_report_ts"] = now
                            self._persist_locked(tk)
            except Exception as e:
                logger.error(f"[process] 无进度重试重投 job_id={RETRY_JOB_ID} 失败: {e}")
                if queued_local == 0:
                    return {"ok": False, "requeued": 0, "missing": missing,
                            "error": f"重投远端失败: {e}"}

        total = queued_local + remote_n
        if total == 0:
            reason = "缺少原始种子 raw_seed" if missing else "无匹配任务"
            return {"ok": False, "requeued": 0, "missing": missing,
                    "error": f"无可重试任务（{reason}）"}
        logger.info(f"[process] 重试 {total} 个：有进度原机重跑 {queued_local}，"
                    f"无进度重投 job_id {remote_n}（缺原始种子 {missing}）")
        return {"ok": True, "requeued": total, "local": queued_local,
                "remote": remote_n, "missing": missing}

    def claim_retries(self, machine: str) -> List[str]:
        """原机 downloader 轮询取回属于本机的待重试原始种子（取出即清空，避免重复）。"""
        if not machine:
            return []
        with self._lock:
            seeds = self._retry_queue.pop(machine, [])
            if seeds:
                try:                              # 取出即从 DB 清除，避免重复重跑
                    self._db.execute(
                        "DELETE FROM retry_queue WHERE machine=?", (machine,))
                    self._db.commit()
                except Exception as e:
                    logger.warning(f"[process] 清除重试队列失败: {e}")
        if seeds:
            logger.info(f"[process] 机器 {machine} 取回 {len(seeds)} 个本机重试种子")
        return list(seeds)

    def reroute_stale_retries(self, max_age: float) -> int:
        """把长时间无人认领的本机重试种子改投远端 job_id，交给任意空闲机器重下。

        “有进度”的重试会绑定到原下载机器的本机队列（复用本地文件），仅由原机的重试轮询
        线程 claim。但若原机已空闲（批次结束、下载进程退出）、失联、或 MACHINE_ID 变化
        （含 IP 变动），这些种子无人 claim，会永久积压在原机队列——即使其它机器空闲也接
        不到。此处兜底：滞留超过 max_age 的条目视为原机不再消费，改投远端 RETRY_JOB_ID，
        由调度分配到任意机器重新下载（放弃本地文件复用，换取不再永久积压）。
        """
        now = time.time()
        with self._lock:
            try:
                cur = self._db.execute(
                    "SELECT id, machine, raw_seed FROM retry_queue WHERE ts < ?",
                    (now - max_age,))
                rows = [(rid, m, r) for rid, m, r in cur.fetchall() if r]
            except Exception as e:
                logger.warning(f"[process] 扫描过期重试队列失败: {e}")
                return 0
            if not rows:
                return 0
            # 先“认领”（从 DB + 内存队列移除），避免原机同时 claim 造成两处重复处理；
            # 若随后远端重投失败再回滚入库。
            try:
                self._db.executemany(
                    "DELETE FROM retry_queue WHERE id=?",
                    [(rid,) for rid, _m, _r in rows])
                self._db.commit()
            except Exception as e:
                logger.warning(f"[process] 认领过期重试队列失败: {e}")
                return 0
            for _rid, m, r in rows:
                q = self._retry_queue.get(m)
                if q:
                    try:
                        q.remove(r)
                    except ValueError:
                        pass
                    if not q:
                        self._retry_queue.pop(m, None)

        seeds = [r for _rid, _m, r in rows]
        try:
            try:
                from .common import XunLiuClient, split_by_size
            except ImportError:
                from common import XunLiuClient, split_by_size
            client = XunLiuClient()
            for chunk in split_by_size(seeds, 40 * 1024 * 1024):
                client.upload_seed(RETRY_JOB_ID, chunk)
        except Exception as e:
            # 重投失败：回滚，保留在原机队列等待下次扫描/原机恢复
            with self._lock:
                for _rid, m, r in rows:
                    self._retry_queue.setdefault(m, []).append(r)
                    try:
                        self._db.execute(
                            "INSERT INTO retry_queue(machine, raw_seed, ts) "
                            "VALUES(?,?,?)", (m, r, now))
                    except Exception:
                        pass
                try:
                    self._db.commit()
                except Exception:
                    pass
            logger.error(f"[process] 过期重试改投远端 job_id={RETRY_JOB_ID} 失败，"
                         f"已回滚保留: {e}")
            return 0

        machines = sorted({m for _rid, m, _r in rows})
        logger.warning(f"[process] {len(rows)} 个重试种子在原机队列滞留超过 "
                       f"{int(max_age)}s（原机空闲/失联/机器名变化: {machines}），"
                       f"已改投远端 job_id={RETRY_JOB_ID} 交由任意空闲机器重新下载")
        return len(rows)

    def remove(self, task_key: str) -> None:
        with self._lock:
            self._status.pop(task_key, None)
            self._desired.pop(task_key, None)
            self._added_ts.pop(task_key, None)
            self._dirty.discard(task_key)
            self._delete_locked(task_key)

    # -------------------------------------------------- 查询（供接口）
    def _passes(self, snap: dict, rid, start: Optional[float],
                end: Optional[float], machine=None, q=None) -> bool:
        if rid not in (None, "", "all"):
            try:
                if int(snap.get("requirement_id") or 0) != int(rid):
                    return False
            except (TypeError, ValueError):
                return False
        ats = float(snap.get("_added_ts", 0) or 0)
        if start is not None and ats < start:
            return False
        if end is not None and ats > end:
            return False
        if machine not in (None, "", "all"):
            if (snap.get("machine") or "") != machine:
                return False
        if q:
            ql = str(q).strip().lower()
            if ql:
                hay = " ".join((
                    str(snap.get("info_hash") or ""),
                    str(snap.get("info_hash_v2") or ""),
                    str(snap.get("name") or ""),
                )).lower()
                if ql not in hay:
                    return False
        return True

    def _light_row(self, key: str, snap: dict, desired: str, now: float) -> dict:
        state = snap.get("state", "")
        offline = (now - snap.get("_report_ts", 0)) > OFFLINE_SECONDS
        return {
            "task_key": key,
            "info_hash": snap.get("info_hash", "") or "",
            "requirement_id": int(snap.get("requirement_id") or 0),
            "name": snap.get("name", "") or "",
            "machine": snap.get("machine", "") or "",
            "state": state,
            "state_label": STATE_LABELS.get(state, state),
            "desired": desired,
            "offline": offline,
            "progress": float(snap.get("progress") or 0),
            "download_rate": float(snap.get("download_rate") or 0),
            "upload_rate": float(snap.get("upload_rate") or 0),
            "num_seeds": int(snap.get("num_seeds") or 0),
            "num_peers": int(snap.get("num_peers") or 0),
            "num_files": int(snap.get("num_files") or 0),
            "completed_files": int(snap.get("completed_files") or 0),
            "total_size": int(snap.get("total_size") or 0),
            "downloaded": int(snap.get("downloaded") or 0),
            "uploaded": int(snap.get("uploaded") or 0),
            "added_ts": float(snap.get("_added_ts", 0) or 0),
            "updated_ts": float(snap.get("_report_ts", 0) or 0),
        }

    def summary(self, rid=None, start: Optional[float] = None,
                end: Optional[float] = None, machine=None, q=None) -> dict:
        counts = {k: 0 for k in GROUP_KEYS}
        total_rate = 0.0
        seeding_count = 0
        seed_up_rate = 0.0
        success_dl_bytes = 0      # 已成功下载的数据总量（下载完成的种子体量，按 info_hash 去重）
        _counted_ih: set = set()  # 去重：多 rid 复用同一 info_hash 只计一次物理下载量
        machines: Dict[str, dict] = {}
        rids: set = set()
        now = time.time()
        with self._lock:
            for key, snap in self._status.items():
                r = int(snap.get("requirement_id") or 0)
                if r:
                    rids.add(r)                       # 候选下拉：全量（不受筛选影响）
                # 机器卡片：按 rid/时间/搜索 过滤，但**不含机器过滤**，这样切换机器时
                # 所有机器卡片仍在、可点击切换。
                if not self._passes(snap, rid, start, end, machine=None, q=q):
                    continue
                mn = snap.get("machine") or "(未知)"
                m = machines.setdefault(mn, {
                    "name": mn, "tasks": 0, "downloading": 0,
                    "dl_rate": 0.0, "ul_rate": 0.0,
                    "disk": None, "cpu": None, "mem": None, "_stat_ts": 0})
                dk = snap.get("disk")
                if dk and (m["disk"] is None or float(dk.get("percent") or 0)
                           > float(m["disk"].get("percent") or 0)):
                    m["disk"] = dk
                # CPU/内存为整机指标：取该机最近一次上报的值
                rts = float(snap.get("_report_ts", 0) or 0)
                if rts >= float(m["_stat_ts"] or 0) and \
                        (snap.get("cpu_percent") is not None
                         or snap.get("mem_percent") is not None):
                    m["cpu"] = snap.get("cpu_percent")
                    m["mem"] = snap.get("mem_percent")
                    m["_stat_ts"] = rts
                desired = self._desired.get(key, DESIRED_RUNNING)
                state = snap.get("state", "")
                offline = (now - snap.get("_report_ts", 0)) > OFFLINE_SECONDS
                cat = _category_of(state, desired, offline)
                # 每机卡片统计（各机独立，不受全局机器筛选影响）：任务数只算“正在下载”，
                # 并汇总该机下载/上传速度。
                m["tasks"] += 1
                if state == ST_DOWNLOADING and not offline:
                    m["downloading"] += 1        # 只算“正在下载”（state=downloading）
                    m["dl_rate"] += float(snap.get("download_rate") or 0)
                if not offline:
                    m["ul_rate"] += float(snap.get("upload_rate") or 0)
                # 全局计数 / 速度 / 流量：叠加机器过滤（选中机器时只统计该机）
                if machine not in (None, "", "all") and mn != machine:
                    continue
                counts[cat] = counts.get(cat, 0) + 1
                # 已成功下载的数据总量：整种已下完的状态才计入（做种/已入库完成/已完成/
                # 存储不足删除——均为已下完再入库）；按 info_hash 去重避免多 rid 复用重复计。
                if state in (ST_COMPLETED, ST_SEEDING, ST_INGESTED_DONE,
                             ST_DELETED_NO_SPACE):
                    ih = str(snap.get("info_hash") or "")
                    if ih and ih not in _counted_ih:
                        _counted_ih.add(ih)
                        success_dl_bytes += int(snap.get("total_size") or 0)
                if state == ST_DOWNLOADING and not offline:
                    total_rate += float(snap.get("download_rate") or 0)
                if state == ST_SEEDING and not offline:
                    seeding_count += 1
                    seed_up_rate += float(snap.get("upload_rate") or 0)
        return {
            "counts": counts, "total_rate": total_rate,
            "seeding_count": seeding_count, "seed_up_rate": seed_up_rate,
            "success_dl_bytes": success_dl_bytes,
            "machines": machines, "rids": sorted(rids),
            "updated": now,
        }

    def list_page(self, category: str, rid=None, start: Optional[float] = None,
                  end: Optional[float] = None, offset: int = 0,
                  limit: int = 20, machine=None, q=None,
                  sort=None, sort_dir: int = -1) -> dict:
        now = time.time()
        items: List[Tuple[float, str, dict, str]] = []
        with self._lock:
            for key, snap in self._status.items():
                if not self._passes(snap, rid, start, end, machine=machine, q=q):
                    continue
                desired = self._desired.get(key, DESIRED_RUNNING)
                offline = (now - snap.get("_report_ts", 0)) > OFFLINE_SECONDS
                cat = _category_of(snap.get("state", ""), desired, offline)
                if category and category != "all" and cat != category:
                    continue
                items.append((float(snap.get("_added_ts", 0) or 0), key, snap, desired))
            # 排序对**全部**过滤结果生效，再分页（保证“加载更多”与已显示的处于同一全局排序）。
            # 默认按添加时间倒序；指定 sort 列（进度/速度/总大小）则按该列排全量。
            sort_key = (str(sort or "").strip())
            colmap = {
                "progress": lambda t: float(t[2].get("progress") or 0),
                "download_rate": lambda t: float(t[2].get("download_rate") or 0),
                "total_size": lambda t: float(t[2].get("total_size") or 0),
                "machine": lambda t: str(t[2].get("machine") or ""),
                "added_ts": lambda t: t[0],
            }
            if sort_key in colmap and sort_key != "added_ts":
                # 主排序键相等时用添加时间倒序做稳定次序
                items.sort(key=lambda x: x[0], reverse=True)
                items.sort(key=colmap[sort_key], reverse=(int(sort_dir) < 0))
            else:
                items.sort(key=lambda x: x[0], reverse=True)
            total = len(items)
            page = items[offset:offset + limit] if limit else items[offset:]
            rows = [self._light_row(key, snap, desired, now)
                    for _ats, key, snap, desired in page]
        return {"rows": rows, "total": total, "offset": offset, "limit": limit}

    def export_rows(self, category: str, rid=None, start: Optional[float] = None,
                    end: Optional[float] = None, machine=None, q=None) -> List[dict]:
        return self.list_page(category, rid, start, end, offset=0, limit=0,
                              machine=machine, q=q)["rows"]

    def export_seed_rows(self, category: str, rid=None,
                         start: Optional[float] = None, end: Optional[float] = None,
                         machine=None, q=None) -> List[Tuple[str, dict]]:
        """导出用：返回 [(raw_seed 原始种子串, status 现有逻辑状态)]，按添加时间倒序。

        raw_seed 为 downloader 上报的原始种子原文（前端“重试”也用它）；status 汇总当前
        逻辑判定的状态（分类/状态文案/进度/机器/速度/错误等），供导出成 JSONL。
        """
        now = time.time()
        out: List[Tuple[float, str, dict]] = []
        with self._lock:
            for key, snap in self._status.items():
                if not self._passes(snap, rid, start, end, machine=machine, q=q):
                    continue
                desired = self._desired.get(key, DESIRED_RUNNING)
                offline = (now - snap.get("_report_ts", 0)) > OFFLINE_SECONDS
                state = snap.get("state", "")
                cat = _category_of(state, desired, offline)
                if category and category != "all" and cat != category:
                    continue
                status = {
                    "task_key": key,
                    "info_hash": snap.get("info_hash", "") or "",
                    "info_hash_v2": snap.get("info_hash_v2", "") or "",
                    "requirement_id": int(snap.get("requirement_id") or 0),
                    "name": snap.get("name", "") or "",
                    "machine": snap.get("machine", "") or "",
                    "state": state,
                    "state_label": STATE_LABELS.get(state, state),
                    "category": cat,
                    "desired": desired,
                    "offline": offline,
                    "progress": float(snap.get("progress") or 0),
                    "download_rate": float(snap.get("download_rate") or 0),
                    "upload_rate": float(snap.get("upload_rate") or 0),
                    "num_seeds": int(snap.get("num_seeds") or 0),
                    "num_peers": int(snap.get("num_peers") or 0),
                    "completed_files": int(snap.get("completed_files") or 0),
                    "num_files": int(snap.get("num_files") or 0),
                    "total_size": int(snap.get("total_size") or 0),
                    "error": snap.get("error", "") or "",
                    "added_ts": float(snap.get("_added_ts", 0) or 0),
                    "updated_ts": float(snap.get("_report_ts", 0) or 0),
                }
                out.append((float(snap.get("_added_ts", 0) or 0),
                            snap.get("raw_seed") or "", status))
            out.sort(key=lambda x: x[0], reverse=True)
        return [(raw, status) for _ats, raw, status in out]

    def get_task(self, task_key: str) -> Optional[dict]:
        now = time.time()
        with self._lock:
            snap = self._status.get(task_key)
            if snap is None:
                return None
            s = dict(snap)
        s["task_key"] = task_key
        s["desired"] = self._desired.get(task_key, DESIRED_RUNNING)
        s["offline"] = (now - snap.get("_report_ts", 0)) > OFFLINE_SECONDS
        state = snap.get("state", "")
        s["state_label"] = STATE_LABELS.get(state, state)
        return s

    def all_status(self) -> Dict[str, dict]:
        """全量状态（兼容旧 /api/status 与控制台展示）。"""
        with self._lock:
            out = {}
            now = time.time()
            for key, snap in self._status.items():
                s = dict(snap)
                s["task_key"] = key
                s["desired"] = self._desired.get(key, DESIRED_RUNNING)
                s["offline"] = (now - snap.get("_report_ts", 0)) > OFFLINE_SECONDS
                state = s.get("state", "")
                s["state_label"] = STATE_LABELS.get(state, state)
                out[key] = s
            return out


STORE = MonitorStore()


# ============================================================================
# HTTP 服务端
# ============================================================================

def _render_dashboard() -> str:
    """看板静态页：数据由前端按分类分页 fetch，无限下拉懒加载，静默增量刷新。"""
    return """<!doctype html>
<html lang="zh-CN">
<head>
<meta charset="utf-8">
<meta name="viewport" content="width=device-width, initial-scale=1">
<title>torrent 下载监控</title>
<style>
:root{
  --bg:#0e1117; --panel:#161b22; --panel2:#1c2230; --border:#2a3040;
  --text:#e6edf3; --muted:#8b949e; --accent:#58a6ff;
  --green:#3fb950; --teal:#39c5cf; --orange:#d29922; --red:#f85149;
  --blue:#58a6ff; --gray:#6e7681;
}
*{box-sizing:border-box}
body{margin:0;font-family:-apple-system,BlinkMacSystemFont,"Segoe UI",
  "PingFang SC","Microsoft YaHei",Helvetica,Arial,sans-serif;
  background:var(--bg);color:var(--text);font-size:14px}
header{position:sticky;top:0;z-index:10;background:linear-gradient(180deg,#121620,#0e1117);
  border-bottom:1px solid var(--border);padding:16px 24px;
  display:flex;align-items:center;gap:24px;flex-wrap:wrap}
header h1{font-size:18px;margin:0;font-weight:600;letter-spacing:.3px}
header h1 .dot{display:inline-block;width:9px;height:9px;border-radius:50%;
  background:var(--green);margin-right:8px;vertical-align:middle;
  box-shadow:0 0 8px var(--green)}
.filters{display:flex;gap:10px;align-items:center;flex-wrap:wrap}
.filters label{color:var(--muted);font-size:12px}
.filters select,.filters input{background:var(--panel);color:var(--text);
  border:1px solid var(--border);border-radius:6px;padding:5px 8px;font:inherit;font-size:12px}
.filters .req{color:var(--red);margin-left:2px}
.filters button{font:inherit;font-size:12px;cursor:pointer;border-radius:6px;
  border:1px solid var(--border);background:var(--panel2);color:var(--text);padding:5px 12px}
.filters button:hover{border-color:var(--accent);color:var(--accent)}
.stats{display:flex;gap:14px;flex-wrap:wrap;margin-left:auto}
.stat{background:var(--panel);border:1px solid var(--border);border-radius:8px;
  padding:8px 14px;min-width:96px}
.stat .k{color:var(--muted);font-size:11px;text-transform:uppercase;letter-spacing:.5px}
.stat .v{font-size:17px;font-weight:600;margin-top:2px}
.stat.clk{cursor:pointer;transition:border-color .15s,background .15s}
.stat.clk:hover{border-color:var(--accent);background:var(--panel2,rgba(88,166,255,.08))}
.mfilter-tag{font-size:12px;color:var(--accent);background:rgba(88,166,255,.14);
  border:1px solid rgba(88,166,255,.4);border-radius:999px;padding:3px 10px;cursor:pointer;
  white-space:nowrap}
.mfilter-tag:hover{background:rgba(248,81,73,.16);color:var(--red);border-color:rgba(248,81,73,.5)}
/* 排序表头 / 复制按钮 / info_hash 展示 / 机器排序条 */
th.sortable{cursor:pointer;user-select:none;white-space:nowrap}
th.sortable:hover{color:var(--accent)}
.copy-btn{font:inherit;font-size:11px;cursor:pointer;border-radius:5px;
  border:1px solid var(--border);background:var(--panel);color:var(--muted);
  padding:1px 8px;margin-left:6px}
.copy-btn:hover{border-color:var(--accent);color:var(--accent)}
.ihbox{display:flex;align-items:center;gap:8px;flex-wrap:wrap;padding:8px 10px;
  background:var(--panel);border:1px solid var(--border);border-radius:8px;margin-bottom:4px}
.ihbox .ihlbl{color:var(--muted);font-size:12px}
.ihbox .ihval{font-family:ui-monospace,Menlo,monospace;font-size:12px;
  user-select:all;word-break:break-all;background:rgba(110,118,129,.12);
  padding:2px 6px;border-radius:5px}
.pctmini{font-size:11px;color:var(--muted);margin-left:6px}
.msortbar{margin-left:12px;font-size:12px;color:var(--muted);font-weight:500}
.msort{font:inherit;font-size:11px;cursor:pointer;border-radius:5px;
  border:1px solid var(--border);background:var(--panel);color:var(--muted);
  padding:1px 8px;margin-left:4px}
.msort:hover{border-color:var(--accent);color:var(--accent)}
.msort.on{border-color:var(--accent);color:var(--accent);background:rgba(88,166,255,.1)}
.wrap{padding:20px 24px}
.group{margin-bottom:22px}
.ghead{font-size:14px;font-weight:600;margin:0 0 10px;display:flex;
  align-items:center;gap:9px;color:var(--muted);cursor:pointer;user-select:none;
  padding:6px 4px;border-radius:8px}
.ghead:hover{background:var(--panel)}
.ghead .cnt{background:var(--panel2);border:1px solid var(--border);
  border-radius:999px;padding:1px 9px;font-size:12px;color:var(--text)}
.ghead .mk{width:9px;height:9px;border-radius:50%}
.ghead .gcaret{color:var(--muted);transition:transform .15s;font-size:12px;width:12px}
.ghead:not(.folded) .gcaret{transform:rotate(90deg);color:var(--accent)}
.ghead .resume-all,.ghead .exp-btn{margin-left:10px;font-size:12px;padding:3px 11px;cursor:pointer;
  border:1px solid var(--border);border-radius:6px;background:var(--panel2);color:var(--text)}
.ghead .resume-all:hover{border-color:var(--green);color:var(--green)}
.ghead .exp-btn{margin-left:8px}
.ghead .exp-btn:hover{border-color:var(--accent);color:var(--accent)}
.mk-dl{background:var(--green)} .mk-pause{background:var(--orange)}
.mk-del{background:var(--red)} .mk-stall{background:#db6d28} .mk-ingest{background:var(--teal)}
.glist{max-height:62vh;overflow:auto;border:1px solid var(--border);border-radius:10px;
  background:var(--panel)}
tr.more td,tr.sentinel td{background:var(--panel);text-align:center;padding:9px;color:var(--muted)}
tr.more button{border:1px dashed var(--border);background:var(--panel);color:var(--fg);
  cursor:pointer;padding:6px 18px;border-radius:8px;font-size:13px}
tr.more button:hover{border-color:var(--muted);border-style:solid}
table{width:100%;border-collapse:separate;border-spacing:0;background:var(--panel)}
.glist>table{border:none;border-radius:0}
thead th{text-align:left;background:var(--panel2);color:var(--muted);
  font-weight:600;font-size:12px;padding:11px 14px;
  text-transform:uppercase;letter-spacing:.4px;border-bottom:1px solid var(--border)}
.glist thead th{position:sticky;top:0;z-index:1}
tbody td{padding:11px 14px;border-bottom:1px solid var(--border);vertical-align:middle}
tbody tr.task{cursor:pointer}
tbody tr.task:hover{background:#1a2130}
tbody tr.task.open{background:#1a2130}
.ih{font-family:ui-monospace,SFMono-Regular,Menlo,monospace;font-weight:600;
  display:flex;align-items:center;gap:8px}
.ih .caret{color:var(--muted);transition:transform .15s;display:inline-block}
tr.task.open .ih .caret{transform:rotate(90deg);color:var(--accent)}
.name{color:var(--muted);font-size:12px;margin-top:3px;max-width:320px;
  overflow:hidden;text-overflow:ellipsis;white-space:nowrap}
.sub{color:var(--muted);font-size:12px;font-family:ui-monospace,SFMono-Regular,Menlo,monospace}
.badge{display:inline-flex;align-items:center;gap:6px;padding:4px 10px;border-radius:999px;
  font-size:12px;font-weight:600;white-space:nowrap}
.badge .b-dot{width:7px;height:7px;border-radius:50%;background:currentColor}
.b-metadata{color:var(--blue);background:rgba(88,166,255,.14)}
.b-metadata .b-dot{animation:pulse 1.1s infinite}
.b-downloading{color:var(--green);background:rgba(63,185,80,.14)}
.b-uploading{color:var(--teal);background:rgba(57,197,207,.14)}
.b-paused{color:var(--orange);background:rgba(210,153,34,.16)}
.b-completed{color:var(--green);background:rgba(63,185,80,.14)}
.b-error{color:var(--red);background:rgba(248,81,73,.16)}
.b-offline{color:var(--gray);background:rgba(110,118,129,.18)}
.b-pending{color:var(--muted);background:rgba(139,148,158,.14)}
.b-stalled{color:#db6d28;background:rgba(219,109,40,.16)}
.b-ingested{color:var(--teal);background:rgba(57,197,207,.14)}
@keyframes pulse{0%,100%{opacity:1}50%{opacity:.25}}
.bar{position:relative;height:8px;border-radius:6px;background:#0d1117;
  border:1px solid var(--border);overflow:hidden;min-width:110px}
.bar > i{position:absolute;left:0;top:0;bottom:0;border-radius:6px;
  background:linear-gradient(90deg,#2a6fdb,#58a6ff)}
.bar.indeterminate > i{width:35%!important;
  background:linear-gradient(90deg,transparent,var(--blue),transparent);
  animation:slide 1.2s infinite}
.bar.done > i{background:linear-gradient(90deg,#2ea043,#3fb950)}
.bar.sm{height:6px;min-width:80px}
@keyframes slide{0%{left:-35%}100%{left:100%}}
.pct{font-size:12px;color:var(--muted);margin-top:4px}
.spd{font-family:ui-monospace,SFMono-Regular,Menlo,monospace;font-size:12px;white-space:nowrap}
.spd .ul{color:var(--muted)}
.btns{display:flex;gap:6px;flex-wrap:wrap}
button.retry:hover{border-color:var(--teal);color:var(--teal)}
button{font:inherit;font-size:12px;cursor:pointer;border-radius:6px;
  border:1px solid var(--border);background:var(--panel2);color:var(--text);
  padding:6px 10px;transition:.15s}
button:hover{border-color:var(--accent);color:var(--accent)}
button.ingest:hover{border-color:var(--teal);color:var(--teal)}
button.danger:hover{border-color:var(--red);color:var(--red)}
button:disabled{opacity:.35;cursor:not-allowed}
.section-div{display:flex;align-items:center;gap:12px;margin:22px 0 6px;color:var(--muted);
  font-size:13px;font-weight:600}
.section-div::before,.section-div::after{content:"";flex:1;height:1px;background:var(--border)}
.empty{text-align:center;color:var(--muted);padding:34px}
.desired{font-size:11px;color:var(--muted);margin-top:3px}
tr.detail > td{background:#12161f;padding:0;border-bottom:1px solid var(--border)}
.det{padding:16px 18px;display:grid;grid-template-columns:1fr 1fr;gap:18px}
.det .full{grid-column:1 / -1}
.panel-title{font-size:12px;font-weight:700;color:var(--muted);
  text-transform:uppercase;letter-spacing:.5px;margin:0 0 8px}
.chips{display:flex;gap:8px;flex-wrap:wrap}
.chip{background:var(--panel2);border:1px solid var(--border);border-radius:8px;
  padding:6px 11px;font-size:12px;display:flex;gap:7px;align-items:center}
.chip b{color:var(--text)} .chip .lab{color:var(--muted)}
.chip.on{border-color:rgba(63,185,80,.5)}
.mini{width:100%;border:1px solid var(--border);border-radius:8px;overflow:hidden;
  background:var(--panel)}
.mini th{font-size:11px;padding:7px 10px;background:var(--panel2);color:var(--muted);
  text-align:left;border-bottom:1px solid var(--border)}
.mini td{font-size:12px;padding:7px 10px;border-bottom:1px solid var(--border)}
.mini tr:last-child td{border-bottom:none}
.mini .u{font-family:ui-monospace,SFMono-Regular,Menlo,monospace;
  max-width:340px;overflow:hidden;text-overflow:ellipsis;white-space:nowrap}
.dot-ok{color:var(--green)} .dot-bad{color:var(--red)} .dot-idle{color:var(--muted)}
.scroll{max-height:240px;overflow:auto}
#toast{position:fixed;bottom:20px;right:20px;background:var(--panel2);
  border:1px solid var(--border);border-radius:8px;padding:10px 16px;
  opacity:0;transform:translateY(8px);transition:.2s;pointer-events:none;z-index:20}
#toast.show{opacity:1;transform:none}
.mtitle{font-size:13px;font-weight:600;color:var(--muted);margin:0 0 10px;
  text-transform:uppercase;letter-spacing:.4px}
.machines{display:flex;gap:12px;flex-wrap:wrap}
.mcard{background:var(--panel);border:1px solid var(--border);border-radius:10px;
  padding:12px 16px;min-width:230px;cursor:pointer;transition:border-color .15s,box-shadow .15s}
.mcard:hover{border-color:var(--accent)}
.mcard.sel{border-color:var(--accent);box-shadow:0 0 0 1px var(--accent) inset}
.mcard.disk-high{border-color:rgba(210,153,34,.6)}
.mcard.disk-danger{border-color:rgba(248,81,73,.75)}
.mcard .mname{font-weight:600;font-family:ui-monospace,Menlo,monospace;
  display:flex;align-items:center;gap:8px;margin-bottom:2px}
.mtag{font-size:11px;padding:2px 8px;border-radius:999px;font-weight:600;margin-left:auto}
.mtag.high{color:var(--orange);background:rgba(210,153,34,.16)}
.mtag.danger{color:var(--red);background:rgba(248,81,73,.16)}
.mbar{position:relative;height:8px;border-radius:6px;background:#0d1117;
  border:1px solid var(--border);overflow:hidden;margin-top:8px}
.mbar>i{position:absolute;left:0;top:0;bottom:0;border-radius:6px;
  background:linear-gradient(90deg,#2ea043,#3fb950)}
.mbar.high>i{background:linear-gradient(90deg,#bb8009,#d29922)}
.mbar.danger>i{background:linear-gradient(90deg,#c93c37,#f85149)}
.mmeta{color:var(--muted);font-size:12px;margin-top:7px;display:flex;
  justify-content:space-between;gap:12px}
.mstat{margin-top:8px;display:flex;flex-direction:column;gap:6px}
.mstat .ms{display:flex;align-items:center;gap:8px}
.mstat .mslab{color:var(--muted);font-size:11px;min-width:74px}
.mstat .mbar{flex:1;height:6px}
.alerts{display:flex;flex-direction:column;gap:8px}
.alert{display:flex;align-items:center;gap:10px;padding:9px 14px;border-radius:8px;
  border:1px solid var(--border);background:var(--panel);font-size:13px}
.alert.high{border-color:rgba(210,153,34,.5);background:rgba(210,153,34,.08)}
.alert.danger{border-color:rgba(248,81,73,.55);background:rgba(248,81,73,.08)}
.alert .lvl{font-weight:600;padding:2px 9px;border-radius:999px;font-size:11px;white-space:nowrap}
.alert.high .lvl{color:var(--orange);background:rgba(210,153,34,.16)}
.alert.danger .lvl{color:var(--red);background:rgba(248,81,73,.16)}
.alert .amac{font-family:ui-monospace,Menlo,monospace;font-weight:600}
.alert .amsg{color:var(--muted)}
/* 消息列表（历史错误 / 告警） */
.ev-head{display:flex;align-items:center;gap:10px;flex-wrap:wrap}
.ev-q{flex:1;min-width:200px;max-width:420px;padding:7px 12px;border-radius:8px;
  border:1px solid var(--border);background:var(--panel);color:var(--fg);font-size:13px}
.ev-q:focus{outline:none;border-color:var(--accent)}
.ev-toggle{display:inline-flex;align-items:center;gap:8px;cursor:pointer;
  padding:8px 14px;border-radius:8px;border:1px solid var(--border);
  background:var(--panel);color:var(--fg);font-size:13px;font-weight:600}
.ev-toggle:hover{border-color:var(--muted)}
.ev-toggle .gcaret{display:inline-block;transition:transform .15s;color:var(--muted)}
.ev-toggle.open .gcaret{transform:rotate(90deg)}
.ev-toggle .ev-cnt{color:var(--muted);font-weight:500}
.evlist{max-height:34vh;overflow:auto;border:1px solid var(--border);
  border-radius:10px;background:var(--panel);margin-top:10px}
.evlist thead th{position:sticky;top:0;z-index:1}
.ev-error{color:var(--red)} .ev-warning{color:var(--orange)} .ev-info{color:var(--muted)}
@media(max-width:900px){.det{grid-template-columns:1fr}}
</style>
</head>
<body>
<header>
  <h1><span class="dot"></span>torrent 下载监控</h1>
  <div class="filters">
    <label>需求ID<span class="req">*</span></label>
    <select id="f-rid"><option value="all">全部</option></select>
    <label>添加时间</label>
    <input type="datetime-local" id="f-start" title="添加起始时间">
    <span style="color:var(--muted)">~</span>
    <input type="datetime-local" id="f-end" title="添加结束时间">
    <input type="text" id="f-q" placeholder="搜索 info_hash / 任务名 / 种子name"
           title="搜索 info_hash 或任务/种子名称" style="min-width:240px">
    <button onclick="clearFilters()">清除</button>
    <span id="f-machine-tag" class="mfilter-tag" style="display:none"></span>
  </div>
  <div class="stats">
    <div class="stat"><div class="k">总下载速度</div><div class="v" id="s-rate">-</div></div>
    <div class="stat"><div class="k">已成功下载数据总量</div><div class="v" id="s-dl-success">-</div></div>
    <div class="stat clk" onclick="jumpGroup('queued')"><div class="k">排队等待</div><div class="v" id="s-queued">-</div></div>
    <div class="stat clk" onclick="jumpGroup('downloading')"><div class="k">下载中</div><div class="v" id="s-dl">-</div></div>
    <div class="stat clk" onclick="jumpGroup('seeding')"><div class="k">做种中</div><div class="v" id="s-seed">-</div></div>
    <div class="stat clk" onclick="jumpGroup('ingested_done')"><div class="k">已全部入库完成</div><div class="v" id="s-ingdone">-</div></div>
    <div class="stat clk" onclick="jumpGroup('paused')"><div class="k">已暂停</div><div class="v" id="s-pause">-</div></div>
    <div class="stat clk" onclick="jumpGroup('stalled')"><div class="k">停滞放弃</div><div class="v" id="s-stall">-</div></div>
    <div class="stat clk" onclick="jumpGroup('extract_failed')"><div class="k">解压/入库失败</div><div class="v" id="s-extract">-</div></div>
    <div class="stat clk" onclick="jumpGroup('force_ingested')"><div class="k">强制入库</div><div class="v" id="s-ingest">-</div></div>
    <div class="stat clk" onclick="jumpGroup('requeued')"><div class="k">重投队列</div><div class="v" id="s-requeued">-</div></div>
    <div class="stat clk" onclick="jumpGroup('other')"><div class="k">其他特殊情况</div><div class="v" id="s-other">-</div></div>
    <div class="stat clk" onclick="jumpGroup('canceled_no_space')"><div class="k">存储不足取消</div><div class="v" id="s-cancel-ns">-</div></div>
    <div class="stat clk" onclick="jumpGroup('deleted_no_space')"><div class="k">存储不足删除</div><div class="v" id="s-del-ns">-</div></div>
    <div class="stat clk" onclick="jumpGroup('deleted')"><div class="k">已删除</div><div class="v" id="s-del">-</div></div>
    <div class="stat"><div class="k">更新时间</div><div class="v" id="s-time">-</div></div>
  </div>
</header>
<div class="wrap" id="machines" style="padding-bottom:0"></div>
<div class="wrap" id="alerts" style="padding-top:8px;padding-bottom:0"></div>
<div class="wrap" id="events" style="padding-top:8px;padding-bottom:0"></div>
<div class="wrap" id="groups"><div class="empty">加载中…</div></div>
<div id="toast"></div>
<script>
const POLL_MS = 2000;
const PAGE = 10;             // 每页/每次点击展示的行数（初始也只展示 10 条）
const DISK_HIGH = 0.8, DISK_PAUSE = 0.9;
const GROUPS = [
  ['queued','排队等待下载','mk-dl'],
  ['downloading','正在下载','mk-dl'],
  ['paused','已暂停','mk-pause'],
  ['stalled','下载停滞放弃','mk-stall'],
  ['extract_failed','解压/入库失败','mk-stall'],
  ['force_ingested','已强制入库','mk-ingest'],
  ['canceled_no_space','因存储不足取消下载','mk-del'],
  ['deleted_no_space','因存储不足已下载入库后删除','mk-del'],
  ['deleted','已删除','mk-del'],
  ['other','其他特殊情况','mk-stall'],
  ['requeued','重投队列','mk-dl'],
  ['seeding','做种中','mk-ingest'],
  ['ingested_done','已全部入库完成','mk-ingest'],
];
const $ = (id) => document.getElementById(id);
const expanded = new Set();       // 展开明细的 info_hash
// 默认全部折叠：一开始进入或刷新页面后每个分类保持关闭，点击标题再展开加载。
const foldedGroups = new Set(GROUPS.map(g=>g[0]));
const detailCache = {};           // info_hash -> 完整明细快照（按需拉取，缓存）
const fileLimits = {};            // info_hash -> 文件明细懒加载上限
const groupState = {};            // key -> {rows, offset, total, done, loading, _sig}
let groupTotals = {};             // 分类计数（来自 summary）
let filter = {rid:'all', start:'', end:'', machine:'', q:''};
const fileSort = {};              // tk -> {col:'rel_path'|'speed'|'progress'|'size', dir:1|-1}
const groupSort = {};             // 分类 key -> {col:'progress'|'download_rate'|..., dir}
let machineSort = {col:'name', dir:1};   // 机器卡片排序
let scrollingUntil = 0;           // 活跃滚动期间跳过刷新，防抖
let fileLoadTs = 0;
const EV_PAGE = 30;               // 消息列表每页/每次下拉加载条数
let evState = {rows:[], offset:0, total:0, done:false, loading:false, _sig:null};
let evOpen = false;               // 消息列表默认收起，点按钮展开
let evQ = '';                     // 消息列表内独立搜索词（info_hash / 名称 / 错误信息）
let _evSearchTimer = null;

function esc(s){return String(s==null?'':s).replace(/[&<>"']/g,
  c=>({'&':'&amp;','<':'&lt;','>':'&gt;','"':'&quot;',"'":'&#39;'}[c]));}
function fmtSize(n){
  n = Number(n)||0; const u=['B','KB','MB','GB','TB']; let i=0;
  while(n>=1024 && i<u.length-1){n/=1024;i++;} return n.toFixed(1)+u[i];
}
function fmtSpeed(n){return fmtSize(n)+'/s';}
function fmtTime(ts){ if(!ts) return '-';
  try{return new Date(ts*1000).toLocaleString('zh-CN',{hour12:false});}catch(e){return '-';}}
function hasSelection(){const s=window.getSelection&&window.getSelection();
  return s && String(s).length>0;}

// ---------------- 排序辅助 ----------------
// 字符串用 zh 本地化比较（中文按拼音首字母、数字自然序）；数值按大小。
function cmpVal(x, y, dir){
  const bothNum = (typeof x==='number' || (!isNaN(parseFloat(x)) && isFinite(x)))
               && (typeof y==='number' || (!isNaN(parseFloat(y)) && isFinite(y)));
  if(bothNum) return dir * ((Number(x)||0) - (Number(y)||0));
  return dir * String(x==null?'':x).localeCompare(String(y==null?'':y),
                                                  'zh-Hans-CN', {numeric:true});
}
function arrowFor(active, dir){ return active ? (dir>0?' ▲':' ▼') : ''; }
function toggleDir(st, col){
  // 同列切换升/降；换列则默认升序（相对路径/字符串）或降序（速度/进度更常看大的）
  if(st && st.col===col){ st.dir = -st.dir; return st; }
  const desc = (col==='speed'||col==='progress'||col==='download_rate'
                ||col==='total_size'||col==='dl_rate'||col==='ul_rate'
                ||col==='downloading'||col==='size');
  return {col, dir: desc ? -1 : 1};
}

function sortFiles(tk, col){
  fileSort[tk] = toggleDir(fileSort[tk], col);
  const k = groupOf(tk); if(k) renderGroup(k, true);
}
function sortGroup(key, col){
  groupSort[key] = toggleDir(groupSort[key], col);
  // 排序改由服务端对**全部任务**排序后分页：重置分页从第 1 页按新排序重新拉取，
  // 之后“加载更多”沿用同一排序，保证未展示的任务也在全局排序里。
  loadPage(key, false);
}
function groupSortParams(key){
  const gs = groupSort[key];
  return gs ? {sort: gs.col, dir: gs.dir} : {};
}
function sortMachines(col){
  machineSort = toggleDir(machineSort, col);
  renderMachines(lastMachines||{});
}

async function copyText(txt, ev){
  try{ await navigator.clipboard.writeText(txt); toast('已复制'); }
  catch(e){
    try{ const t=document.createElement('textarea'); t.value=txt;
      document.body.appendChild(t); t.select(); document.execCommand('copy');
      t.remove(); toast('已复制'); }catch(_){ toast('复制失败'); }
  }
  if(ev) ev.stopPropagation();
}

function fq(extra){
  const p = new URLSearchParams();
  if(filter.rid && filter.rid!=='all') p.set('requirement_id', filter.rid);
  if(filter.start) p.set('start', filter.start);
  if(filter.end) p.set('end', filter.end);
  if(filter.machine) p.set('machine', filter.machine);
  if(filter.q) p.set('q', filter.q);
  for(const k in (extra||{})) p.set(k, extra[k]);
  return p.toString();
}

function resetGroups(){
  for(const [key] of GROUPS)
    groupState[key] = {rows:[], offset:0, total:0, done:false, loading:false, _sig:null};
}

function badgeOf(s){
  if(s.offline) return {cls:'b-offline', txt:'离线'};
  const st = s.state||''; const label = s.state_label || st;
  const map={metadata:'b-metadata',checking:'b-metadata',downloading:'b-downloading',
    uploading:'b-uploading',
    paused:'b-paused','已完成':'b-completed',error:'b-error',pending:'b-pending',
    seeding:'b-uploading',
    deleted:'b-offline',skipped:'b-pending',stalled:'b-stalled',
    extract_failed:'b-error',force_ingested:'b-ingested',
    canceled_no_space:'b-error',deleted_no_space:'b-offline',requeued:'b-pending',
    ingested_done:'b-completed',queued:'b-pending'};
  return {cls:(map[st]||'b-pending'), txt:label};
}

function groupOf(tk){
  for(const [key] of GROUPS)
    if(groupState[key].rows.some(r=>r.task_key===tk)) return key;
  return null;
}

// ---------------- 骨架（只建一次，之后仅更新各分类 glist / 计数） ----------------
const RETRY_ALL_GROUPS = ['stalled','extract_failed','canceled_no_space','deleted','other'];
// 底部独立分区标题（把“重投队列”和“做种列表”彼此分开、也和上面的下载类分组分开）
const SECTION_TITLES = {
  requeued: '重投队列（已重投，等待重新下载）',
  seeding: '做种列表（下载完成后持续上传/分享，不参与重试）',
  ingested_done: '已全部入库完成（做种结束、进程已释放，数据已入库 COS、本地已清）',
};
function buildSkeleton(){
  const html = GROUPS.map(([key,label,mk])=>{
    let act = '';
    if(key==='paused')
      act = `<button class="resume-all" onclick="event.stopPropagation();resumeAll('${key}')">全部继续</button>`;
    else if(RETRY_ALL_GROUPS.includes(key))   // 停滞放弃/存储不足取消/已删除/其他：全部重试
      act = `<button class="resume-all" onclick="event.stopPropagation();retryAll('${key}')">全部重试</button>`;
    const divider = SECTION_TITLES[key]
      ? `<div class="section-div"><span>${SECTION_TITLES[key]}</span></div>` : '';
    return divider + `<div class="group" id="grp-${key}">
      <h2 class="ghead ${foldedGroups.has(key)?'folded':''}" onclick="toggleGroup('${key}')">
        <span class="gcaret">&#9656;</span><span class="mk ${mk}"></span>${label}
        <span class="cnt">${groupTotals[key]||0}</span>${act}
        <button class="exp-btn" onclick="event.stopPropagation();exportGroup('${key}')">导出</button>
      </h2>
      <div class="glist" onscroll="onGroupScroll(event,'${key}')"></div>
    </div>`;
  }).join('');
  $('groups').innerHTML = html;
  for(const [key] of GROUPS) renderGroup(key, true);
}

function taskRow(r){
  const ih = r.info_hash||'';
  const tk = r.task_key;
  const b = badgeOf(r);
  const prog = Math.max(0, Math.min(1, Number(r.progress)||0));
  const pct = (prog*100).toFixed(1);
  const isMeta = (r.state==='metadata') && !r.offline;
  const isDone = (r.state==='已完成') || prog>=1;
  const barCls = 'bar' + (isMeta?' indeterminate':'') + (isDone?' done':'');
  const barInner = isMeta ? '<i></i>' : `<i style="width:${pct}%"></i>`;
  const completed = r.completed_files!=null? r.completed_files : 0;
  const nfiles = r.num_files!=null? r.num_files : 0;
  const dl = (r.state==='downloading' && !r.offline)? (Number(r.download_rate)||0):0;
  const up = (!r.offline)? (Number(r.upload_rate)||0):0;
  const desired = r.desired||'running';
  const open = expanded.has(tk);
  const pd = (desired==='paused')?'disabled':'';
  const sd = (desired==='running')?'disabled':'';
  const del = (desired==='deleted' || r.state==='deleted');
  // 重试按钮：除“正在下载中”“正常暂停”“做种中(在上传)”“已在重投队列”外均可重试
  const inProgress = ['downloading','metadata','checking','uploading','pending'].includes(r.state) && !r.offline;
  const isPaused = (desired==='paused' || r.state==='paused');
  const isSeeding = (r.state==='seeding');   // 做种=正在上传/分享，不重投
  const isRequeued = (r.state==='requeued'); // 已重投、等待重新下载，避免重复投递
  const canRetry = !inProgress && !isPaused && !isSeeding && !isRequeued;
  const spd = `<div class="dl">↓${dl>0?fmtSpeed(dl):'-'}</div>`
            + `<div class="ul sub">↑${up>0?fmtSpeed(up):'-'}</div>`;
  const main = `<tr class="task ${open?'open':''}" onclick="toggle('${esc(tk)}')">
    <td><div class="ih"><span class="caret">&#9656;</span>${esc(ih.slice(0,32))||'(未知)'}</div>
        <div class="name" title="${esc(r.name)}">${esc(r.name)||''}</div>
        <div class="desired">需求 ${esc(r.requirement_id)} · 添加 ${esc(fmtTime(r.added_ts))}</div></td>
    <td class="sub">${esc(r.machine)}</td>
    <td><span class="badge ${b.cls}"><span class="b-dot"></span>${esc(b.txt)}</span>
        ${desired!=='running'?`<div class="desired">期望: ${esc(desired)}</div>`:''}</td>
    <td style="min-width:140px"><div class="${barCls}">${barInner}</div>
        <div class="pct">${isMeta?'解析元数据…':pct+'%'}</div></td>
    <td class="spd">${spd}</td>
    <td class="sub">${fmtSize(r.total_size)}</td>
    <td class="sub">${Number(r.num_seeds)||0} / ${Number(r.num_peers)||0}</td>
    <td class="sub">${completed}/${nfiles}</td>
    <td onclick="event.stopPropagation()"><div class="btns">
      ${del?'':(isSeeding
        ? `<button class="ingest" onclick="ctl('${esc(tk)}','force_ingest')">强制入库(结束做种)</button>`
        : `<button ${pd} onclick="ctl('${esc(tk)}','pause')">暂停</button>
      <button ${sd} onclick="ctl('${esc(tk)}','start')">开始</button>
      <button class="ingest" onclick="ctl('${esc(tk)}','force_ingest')">强制入库</button>`)}
      ${canRetry?`<button class="retry" onclick="retry('${esc(tk)}')">重试</button>`:''}
      <button class="danger" onclick="ctl('${esc(tk)}','delete')">删除</button>
    </div></td>
  </tr>`;
  const detail = open ? `<tr class="detail"><td colspan="9">${detailHtml(tk)}</td></tr>` : '';
  return main + detail;
}

function detailHtml(tk){
  const s = detailCache[tk];
  if(!s) return '<div class="det"><div class="full sub">加载明细中…</div></div>';
  const bt = s.bt||{}; const ps = bt.peer_sources||{};
  const trackers = bt.trackers||[]; const peers = bt.peers||[]; const files = s.files||[];
  const srcDefs=[['tracker','Tracker'],['dht','DHT'],['pex','PEX'],
    ['lsd','本地(LSD)'],['incoming','入站'],['resume','恢复']];
  const chips = srcDefs.map(([k,lab])=>{
    const v=Number(ps[k])||0;
    return `<div class="chip ${v>0?'on':''}"><span class="lab">${lab}</span><b>${v}</b></div>`;
  }).join('');
  const trRows = trackers.length ? trackers.map(t=>{
    const ok=t.working; const err=t.last_error||t.message||'';
    const dot = ok?'<span class="dot-ok">&#9679;</span>'
              : (t.fails>0||t.last_error)?'<span class="dot-bad">&#9679;</span>'
              : '<span class="dot-idle">&#9679;</span>';
    const stTxt = ok?'工作中':(t.fails>0?('失败'+t.fails+'次'):(err?'异常':'待通告'));
    return `<tr><td>${dot}</td><td class="u" title="${esc(t.url)}">${esc(t.url)}</td>
      <td>T${Number(t.tier)||0}</td><td>${esc(stTxt)}</td>
      <td class="sub" title="${esc(err)}">${esc(err.slice(0,40))}</td></tr>`;
  }).join('') : '<tr><td colspan="5" class="sub">暂无 tracker 数据</td></tr>';
  const peerRows = peers.length ? peers.slice(0,40).map(p=>{
    const prog=((Number(p.progress)||0)*100).toFixed(0);
    return `<tr><td class="u">${esc(p.ip)}</td><td class="sub">${esc((p.source||[]).join(','))}</td>
      <td class="u sub" title="${esc(p.from||'')}">${esc(p.from||'-')}</td>
      <td class="sub">${esc(p.client)}</td><td class="spd">${fmtSpeed(p.down)}</td>
      <td>${prog}%</td></tr>`;
  }).join('') : '<tr><td colspan="6" class="sub">暂无已连接 peer</td></tr>';
  // 文件列表排序（点击表头）：相对路径按拼音/首字母，大小/速度/进度按数值
  const fs = fileSort[tk];
  let filesView = files;
  if(fs){
    const kmap = {rel_path:f=>f.rel_path||'', size:f=>Number(f.size)||0,
                  speed:f=>Number(f.speed)||0, progress:f=>Number(f.progress)||0};
    const getv = kmap[fs.col] || kmap.rel_path;
    filesView = files.slice().sort((a,b)=>cmpVal(getv(a), getv(b), fs.dir));
  }
  const flimit = fileLimits[tk] || PAGE;
  const shownFiles = filesView.slice(0, flimit);
  const fileRows = shownFiles.length ? shownFiles.map(f=>{
    const fp=Math.max(0,Math.min(1,Number(f.progress)||0));
    const fpct=(fp*100).toFixed(1);
    const done=(f.status==='已完成')||fp>=1;
    const st=f.status||'';
    const spd=Number(f.speed)||0;
    return `<tr>
      <td class="u" title="${esc(f.rel_path)}">${esc(f.rel_path)}</td>
      <td class="sub">${esc(f.modality||'')}</td>
      <td class="sub">${fmtSize(f.size)}</td>
      <td class="spd">${spd>0?fmtSpeed(spd):'-'}</td>
      <td style="min-width:120px"><div class="bar sm ${done?'done':''}">
        <i style="width:${fpct}%"></i></div><span class="pctmini">${fpct}%</span></td>
      <td class="sub">${esc(st)}${f.error?(' · '+esc(String(f.error).slice(0,40))):''}</td>
    </tr>`;
  }).join('') : '<tr><td colspan="6" class="sub">暂无文件明细</td></tr>';
  const moreFiles = files.length > flimit
      ? `<tr class="more"><td colspan="6">下拉加载更多… 已显示 ${flimit}/${files.length}</td></tr>`
      : '';
  const fSortTh = (label, col, extra) => {
    const active = fs && fs.col===col;
    return `<th class="sortable${extra||''}" onclick="event.stopPropagation();sortFiles('${esc(tk)}','${col}')">`
         + `${label}${arrowFor(active, fs&&fs.dir)}</th>`;
  };
  // 展开详情里再显示一次完整 info_hash（前端列表里显示不全、也不便复制）
  const ihFull = s.info_hash || (tk.indexOf(':')>=0 ? tk.slice(tk.indexOf(':')+1) : tk);
  const ihv2 = s.info_hash_v2 || '';
  const ihBlock = `<div class="full ihbox">
    <span class="ihlbl">info_hash</span>
    <code class="ihval">${esc(ihFull)}</code>
    <button class="copy-btn" onclick="copyText('${esc(ihFull)}',event)">复制</button>
    ${ihv2?`<span class="ihlbl">v2</span><code class="ihval">${esc(ihv2)}</code>`
         + `<button class="copy-btn" onclick="copyText('${esc(ihv2)}',event)">复制</button>`:''}
  </div>`;
  return `<div class="det">
    ${ihBlock}
    <div>
      <p class="panel-title">连接来源（peer 从哪里来）</p>
      <div class="chips">${chips}</div>
    </div>
    <div>
      <p class="panel-title">tracker（建立连接 / 通告状态）</p>
      <div class="scroll" data-sk="tr-${esc(tk)}"><table class="mini"><thead><tr>
        <th></th><th>Tracker</th><th>层</th><th>状态</th><th>信息</th></tr></thead>
        <tbody>${trRows}</tbody></table></div>
    </div>
    <div class="full">
      <p class="panel-title">已连接 peer（${peers.length}）</p>
      <div class="scroll" data-sk="pe-${esc(tk)}"><table class="mini"><thead><tr>
        <th>地址</th><th>来源</th><th>来自(Tracker)</th><th>客户端</th><th>下行</th><th>进度</th></tr></thead>
        <tbody>${peerRows}</tbody></table></div>
    </div>
    <div class="full">
      <p class="panel-title">文件下载进度（${files.length}）</p>
      <div class="scroll" data-sk="fl-${esc(tk)}"
           onscroll="onFileScroll(event,'${esc(tk)}',${files.length})">
        <table class="mini"><thead><tr>
        ${fSortTh('相对路径','rel_path')}<th>模态</th>${fSortTh('大小','size')}
        ${fSortTh('速度','speed')}${fSortTh('进度','progress')}<th>状态</th></tr></thead>
        <tbody>${fileRows}${moreFiles}</tbody></table></div>
    </div>
  </div>`;
}

// 分类内容签名：无变化则跳过重绘（静默刷新防卡顿的核心）
function groupSig(key){
  const g = groupState[key];
  const gs = groupSort[key];
  let s = key+'|'+g.total+'|'+g.done+'|'+foldedGroups.has(key)
        + '|'+(gs?gs.col+gs.dir:'-');
  for(const r of g.rows){
    const tk = r.task_key;
    s += '#'+tk+':'+r.state+':'+(Number(r.progress)||0).toFixed(3)
       +':'+Math.round(Number(r.download_rate)||0)+':'+r.desired+':'+(r.offline?1:0)
       +':'+r.completed_files+'/'+r.num_files+':'+(expanded.has(tk)?1:0);
    if(expanded.has(tk))
      s += 'D'+((detailCache[tk]||{})._report_ts||0)+'f'+(fileLimits[tk]||PAGE)
         + 'S'+((fileSort[tk]&&fileSort[tk].col+fileSort[tk].dir)||'-');
  }
  return s;
}

function renderGroup(key, force){
  const g = groupState[key];
  const cont = document.querySelector('#grp-'+key+' .glist');
  if(!cont) return;
  const sig = groupSig(key);
  if(!force && g._sig===sig) return;   // 无变化，跳过 DOM 重建
  g._sig = sig;
  if(foldedGroups.has(key)){ cont.style.display='none'; return; }
  cont.style.display='';
  // 记录滚动位置（列表容器 + 内部明细滚动容器），重绘后恢复
  const listTop = cont.scrollTop;
  const scrollMap = {};
  cont.querySelectorAll('[data-sk]').forEach(el=>scrollMap[el.getAttribute('data-sk')]=el.scrollTop);
  let body;
  if(!g.rows.length){
    body = `<tr><td class="empty" colspan="9">无</td></tr>`;
  } else {
    // 排序已由服务端对全部任务完成（见 loadPage 的 sort 参数），此处按返回顺序直接渲染。
    body = g.rows.map(r=>taskRow(r)).join('');
    if(!g.done)
      body += `<tr class="more"><td colspan="9">`
            + `<button onclick="loadPage('${key}',true)">点击展示更多（已显示 ${g.rows.length}/${g.total}）</button>`
            + `</td></tr>`;
    else if(g.rows.length>PAGE)
      body += `<tr class="sentinel"><td colspan="9">已全部加载 ${g.rows.length} 条</td></tr>`;
  }
  const gs2 = groupSort[key];
  const gTh = (label, col) => `<th class="sortable" onclick="sortGroup('${key}','${col}')">`
    + `${label}${arrowFor(gs2&&gs2.col===col, gs2&&gs2.dir)}</th>`;
  cont.innerHTML = `<table><thead><tr>
    <th>info_hash</th>${gTh('机器','machine')}<th>状态</th>${gTh('进度','progress')}
    ${gTh('速度(↓/↑)','download_rate')}${gTh('总大小','total_size')}
    <th>做种/连接</th><th>文件</th><th>操作</th>
  </tr></thead><tbody>${body}</tbody></table>`;
  cont.scrollTop = listTop;
  cont.querySelectorAll('[data-sk]').forEach(el=>{
    const v = scrollMap[el.getAttribute('data-sk')]; if(v!=null) el.scrollTop=v;});
}

// ---------------- 分页加载 / 静默刷新 ----------------
async function loadPage(key, append){
  const g = groupState[key];
  if(g.loading || (append && g.done)) return;
  g.loading = true;
  try{
    const off = append ? g.offset : 0;
    const r = await fetch('/api/list?'+fq(Object.assign({category:key, offset:off, limit:PAGE}, groupSortParams(key))),{cache:'no-store'});
    const d = await r.json();
    g.total = d.total;
    g.rows = append ? g.rows.concat(d.rows||[]) : (d.rows||[]);
    g.offset = g.rows.length;
    g.done = g.rows.length >= d.total;
    renderGroup(key, true);
  }catch(e){}
  finally{ g.loading = false; }
}

async function refreshGroup(key){
  const g = groupState[key];
  if(g.loading || g.rows.length===0) return;
  try{
    const cnt = g.rows.length;
    const r = await fetch('/api/list?'+fq(Object.assign({category:key, offset:0, limit:cnt}, groupSortParams(key))),{cache:'no-store'});
    const d = await r.json();
    g.total = d.total;
    g.rows = d.rows||[];
    g.offset = g.rows.length;
    g.done = g.rows.length >= d.total;
    renderGroup(key);   // 由签名判断是否真的需要重绘
  }catch(e){}
}

async function pollSummary(){
  try{
    const r = await fetch('/api/summary?'+fq(),{cache:'no-store'});
    const d = await r.json();
    applySummary(d);
  }catch(e){}
}

function applySummary(d){
  groupTotals = d.counts||{};
  for(const [key] of GROUPS){
    const el = document.querySelector('#grp-'+key+' .cnt');
    if(el) el.textContent = groupTotals[key]||0;
  }
  renderMachines(d.machines||{});
  renderAlerts(d.machines||{});
  populateRids(d.rids||[]);
  $('s-rate').textContent = fmtSpeed(d.total_rate||0);
  $('s-dl-success').textContent = fmtSize(d.success_dl_bytes||0);
  $('s-queued').textContent = groupTotals.queued||0;
  $('s-dl').textContent = groupTotals.downloading||0;
  $('s-seed').textContent = (d.seeding_count||0)
    + ((d.seed_up_rate>0)? (' · ↑'+fmtSpeed(d.seed_up_rate)) : '');
  $('s-ingdone').textContent = groupTotals.ingested_done||0;
  $('s-pause').textContent = groupTotals.paused||0;
  $('s-stall').textContent = groupTotals.stalled||0;
  $('s-extract').textContent = groupTotals.extract_failed||0;
  $('s-ingest').textContent = groupTotals.force_ingested||0;
  $('s-requeued').textContent = groupTotals.requeued||0;
  $('s-other').textContent = groupTotals.other||0;
  $('s-cancel-ns').textContent = groupTotals.canceled_no_space||0;
  $('s-del-ns').textContent = groupTotals.deleted_no_space||0;
  $('s-del').textContent = groupTotals.deleted||0;
  $('s-time').textContent = new Date().toLocaleTimeString('zh-CN',{hour12:false});
}

let ridOptSig = '';
function populateRids(rids){
  const sig = rids.join(',');
  if(sig===ridOptSig) return; ridOptSig = sig;
  const sel = $('f-rid'); const cur = sel.value||'all';
  sel.innerHTML = '<option value="all">全部</option>'
    + rids.map(r=>`<option value="${r}">${r}</option>`).join('');
  sel.value = (cur==='all' || rids.map(String).includes(cur)) ? cur : 'all';
}

async function poll(){
  await pollSummary();
  const busy = (Date.now()<scrollingUntil) || hasSelection();
  for(const [key] of GROUPS){
    if(foldedGroups.has(key)) continue;
    const g = groupState[key];
    if(g.rows.length===0){ if((groupTotals[key]||0)>0) await loadPage(key,false); }
    else if(!busy){ await refreshGroup(key); }
  }
  if(!busy && expanded.size){
    for(const tk of Array.from(expanded)) await loadDetail(tk, true);
  }
  if(evOpen){
    if(evState.rows.length===0){ await loadEvents(false); }
    else if(!busy){ await refreshEvents(); }
  } else if(!busy){
    await refreshEventsCount();   // 收起时仅刷新按钮上的计数徽标
  }
}

async function refreshEventsCount(){
  try{
    const r = await fetch('/api/events?'+evFq({offset:0, limit:1}),{cache:'no-store'});
    const d = await r.json();
    evState.total = d.total||0;
    const cntEl = $('ev-cnt');
    if(cntEl) cntEl.textContent = evState.total ? ('· ' + evState.total + ' 条') : '';
  }catch(e){}
}

// ---------------- 交互（均以 task_key 为键） ----------------
async function loadDetail(tk, rerender){
  try{
    const r = await fetch('/api/task?'+new URLSearchParams({task_key:tk}),{cache:'no-store'});
    if(r.ok){ const d = await r.json(); if(d.task) detailCache[tk] = d.task; }
  }catch(e){}
  if(rerender){ const k = groupOf(tk); if(k) renderGroup(k, true); }
}

async function toggle(tk){
  const k = groupOf(tk);
  if(expanded.has(tk)){ expanded.delete(tk); if(k) renderGroup(k, true); return; }
  expanded.add(tk);
  if(k) renderGroup(k, true);         // 先显示“加载明细中…”
  if(!detailCache[tk]) await loadDetail(tk, true);
}

function toggleGroup(key){
  const head = document.querySelector('#grp-'+key+' .ghead');
  if(foldedGroups.has(key)){
    foldedGroups.delete(key); head.classList.remove('folded');
    if(groupState[key].rows.length===0) loadPage(key,false); else renderGroup(key,true);
  } else {
    foldedGroups.add(key); head.classList.add('folded'); renderGroup(key,true);
  }
}

function onGroupScroll(e, key){
  scrollingUntil = Date.now() + 500;   // 活跃滚动，短暂跳过刷新（改为“点击展示更多”，不再滚动自动加载）
}

function onFileScroll(e, tk, total){
  const el = e.target;
  scrollingUntil = Date.now() + 500;
  if(el.scrollTop + el.clientHeight < el.scrollHeight - 24) return;
  const cur = fileLimits[tk] || PAGE;
  if(cur >= total) return;
  const now = Date.now();
  if(now - fileLoadTs < 250) return;
  fileLoadTs = now;
  fileLimits[tk] = cur + PAGE;
  const k = groupOf(tk); if(k) renderGroup(k, true);
}

async function resumeAll(key){
  toast('正在获取任务列表…');
  let keys = [];
  try{
    const r = await fetch('/api/list?'+fq({category:key, offset:0, limit:1000000}),{cache:'no-store'});
    const d = await r.json(); keys = (d.rows||[]).map(x=>x.task_key);
  }catch(e){}
  if(!keys.length){ toast('无可继续任务'); return; }
  toast('正在继续 '+keys.length+' 个任务…');
  for(const tk of keys){
    try{ await fetch('/api/control',{method:'POST',
      headers:{'Content-Type':'application/json'},
      body:JSON.stringify({task_key:tk,action:'start'})}); }catch(e){}
  }
  toast('已发送“全部继续”'); poll();
}

async function retry(tk){
  try{
    const r = await fetch('/api/retry',{method:'POST',
      headers:{'Content-Type':'application/json'},
      body:JSON.stringify({task_key:tk})});
    const d = await r.json();
    toast(d.ok?('已重投到下载队列'):('重试失败：'+(d.error||''))); poll();
  }catch(e){ toast('重试失败'); }
}

async function retryAll(key){
  toast('正在获取任务列表…');
  let keys = [];
  try{
    const r = await fetch('/api/list?'+fq({category:key, offset:0, limit:1000000}),{cache:'no-store'});
    const d = await r.json(); keys = (d.rows||[]).map(x=>x.task_key);
  }catch(e){}
  if(!keys.length){ toast('无可重试任务'); return; }
  toast('正在重投 '+keys.length+' 个任务…');
  try{
    const r = await fetch('/api/retry',{method:'POST',
      headers:{'Content-Type':'application/json'},
      body:JSON.stringify({task_keys:keys})});
    const d = await r.json();
    toast(d.ok?('已重投 '+(d.requeued||0)+' 个'+(d.missing?('，跳过 '+d.missing+' 个'):'')):('重试失败：'+(d.error||'')));
  }catch(e){ toast('重试失败'); }
  poll();
}

function exportGroup(key){
  const url = '/api/export?'+fq({category:key});
  const a = document.createElement('a');
  a.href = url; a.download = ''; document.body.appendChild(a); a.click(); a.remove();
  const label = (GROUPS.find(g=>g[0]===key)||[])[1]||key;
  toast('正在导出「'+label+'」（当前筛选）');
}

function clearFilters(){
  $('f-rid').value='all'; $('f-start').value=''; $('f-end').value='';
  $('f-q').value=''; filter.machine='';
  onFilterChange();
}
function readFilters(){
  filter.rid = $('f-rid').value||'all';
  const sv = $('f-start').value, ev = $('f-end').value;
  filter.start = sv ? Math.floor(new Date(sv).getTime()/1000) : '';
  filter.end   = ev ? Math.floor(new Date(ev).getTime()/1000) : '';
  filter.q = ($('f-q').value||'').trim();
}
function onFilterChange(){
  readFilters(); resetGroups(); buildSkeleton();
  evState = {rows:[], offset:0, total:0, done:false, loading:false, _sig:null};
  // 搜索时自动展开消息列表，方便直接看该种子的报错/告警
  evOpen = !!filter.q || evOpen;
  initEvents();
  if(evOpen){ const b=$('ev-toggle'); if(b) b.classList.add('open');
    const l=$('ev-list'); if(l) l.style.display=''; }
  poll();
}

// 点击机器卡片：只看该机任务（再次点击同一机器或点标签则取消）
function selectMachine(name){
  filter.machine = (filter.machine===name) ? '' : name;
  resetGroups(); buildSkeleton(); poll();
}

// 点击顶部统计卡片：展开并滚动到对应分组
function jumpGroup(key){
  foldedGroups.delete(key);
  const el = document.querySelector('#grp-'+key);
  if(el){
    const head = el.querySelector('.ghead'); if(head) head.classList.remove('folded');
    if(groupState[key] && groupState[key].rows.length===0) loadPage(key,false);
    else renderGroup(key,true);
    el.scrollIntoView({behavior:'smooth', block:'start'});
  }
}

// ---------------- 消息列表（历史错误 / 告警，按 rid+时间筛选、无限下拉） ----------------
function initEvents(){
  $('events').innerHTML =
      '<div class="ev-head">'
    + '<button id="ev-toggle" class="ev-toggle" onclick="toggleEvents()">'
    + '<span class="gcaret">&#9656;</span>消息列表（历史错误 / 告警）'
    + '<span class="ev-cnt" id="ev-cnt"></span></button>'
    + '<input id="ev-q" class="ev-q" placeholder="搜索 info_hash / 名称 / 错误信息…" '
    + 'oninput="onEvSearch(this.value)">'
    + '</div>'
    + '<div class="evlist" id="ev-list" style="display:none" '
    + 'onscroll="onEventsScroll(event)"></div>';
  renderEvents(true);
}

function onEvSearch(v){
  evQ = (v || '').trim();
  if(!evOpen){                       // 搜索时自动展开消息列表
    evOpen = true;
    const b = $('ev-toggle'); if(b) b.classList.add('open');
    const l = $('ev-list'); if(l) l.style.display = '';
  }
  clearTimeout(_evSearchTimer);
  _evSearchTimer = setTimeout(()=>loadEvents(false), 300);   // 防抖
}

function evFq(extra){
  // 消息列表专用查询：有独立搜索词 evQ 时覆盖全局 q，否则沿用顶部筛选
  const e = Object.assign({}, extra || {});
  if(evQ) e.q = evQ;
  return fq(e);
}

function toggleEvents(){
  evOpen = !evOpen;
  const btn = $('ev-toggle'); const list = $('ev-list');
  if(btn) btn.classList.toggle('open', evOpen);
  if(list) list.style.display = evOpen ? '' : 'none';
  if(evOpen && evState.rows.length===0) loadEvents(false);
  else if(evOpen) renderEvents(true);
}

function renderEvents(force){
  const g = evState;
  const cntEl = $('ev-cnt');
  if(cntEl) cntEl.textContent = g.total ? ('· ' + g.total + ' 条') : '';
  const cont = $('ev-list');
  if(!cont || !evOpen) return;
  const sig = 'e'+g.total+'|'+g.done+'|'
    + g.rows.map(r=>r.ts+':'+r.level+':'+r.info_hash).join(',');
  if(!force && g._sig===sig) return;
  g._sig = sig;
  const top = cont.scrollTop;
  let body;
  if(!g.rows.length){
    body = '<tr><td class="sub" colspan="6" style="text-align:center;padding:16px">'
         + '暂无错误 / 告警消息</td></tr>';
  } else {
    body = g.rows.map(r=>{
      const lv = r.level||'info';
      const cls = lv==='error'?'ev-error':(lv==='warning'?'ev-warning':'ev-info');
      const lvtxt = lv==='error'?'错误':(lv==='warning'?'告警':'信息');
      return `<tr>
        <td class="sub" style="white-space:nowrap">${esc(fmtTime(r.ts))}</td>
        <td><b class="${cls}">${lvtxt}</b></td>
        <td class="sub">需求 ${esc(r.requirement_id)}</td>
        <td class="sub">${esc(r.machine)}</td>
        <td class="u sub" title="${esc(r.info_hash)}">${esc((r.info_hash||'').slice(0,16))}</td>
        <td>${esc(r.message)}</td>
      </tr>`;
    }).join('');
    if(!g.done)
      body += `<tr class="sentinel"><td colspan="6">下拉加载更多… 已显示 ${g.rows.length}/${g.total}</td></tr>`;
  }
  cont.innerHTML = `<table class="mini"><thead><tr>
    <th>时间</th><th>级别</th><th>需求</th><th>机器</th><th>info_hash</th><th>消息</th>
  </tr></thead><tbody>${body}</tbody></table>`;
  cont.scrollTop = top;
}

async function loadEvents(append){
  const g = evState;
  if(g.loading || (append && g.done)) return;
  g.loading = true;
  try{
    const off = append ? g.offset : 0;
    const r = await fetch('/api/events?'+evFq({offset:off, limit:EV_PAGE}),{cache:'no-store'});
    const d = await r.json();
    g.total = d.total||0;
    g.rows = append ? g.rows.concat(d.rows||[]) : (d.rows||[]);
    g.offset = g.rows.length;
    g.done = g.rows.length >= g.total;
    renderEvents(true);
  }catch(e){}
  finally{ g.loading = false; }
}

async function refreshEvents(){
  const g = evState;
  if(g.loading) return;
  try{
    const cnt = Math.max(g.rows.length, EV_PAGE);
    const r = await fetch('/api/events?'+evFq({offset:0, limit:cnt}),{cache:'no-store'});
    const d = await r.json();
    g.total = d.total||0;
    g.rows = d.rows||[];
    g.offset = g.rows.length;
    g.done = g.rows.length >= g.total;
    renderEvents();
  }catch(e){}
}

function onEventsScroll(e){
  scrollingUntil = Date.now() + 500;
  const el = e.target;
  if(el.scrollTop + el.clientHeight < el.scrollHeight - 40) return;
  if(evState.done || evState.loading) return;
  loadEvents(true);
}

// ---------------- 机器磁盘 / 告警（来自 summary） ----------------
let lastMachines = {};
function renderMachines(machines){
  lastMachines = machines || {};
  const list = Object.values(machines);
  if(!list.length){ $('machines').innerHTML=''; return; }
  // 机器卡片排序：名称 / 下载速度 / 下载中任务数（点击排序按钮切换升降）
  const ms = machineSort || {col:'name', dir:1};
  const mkey = {name:m=>m.name||'', dl_rate:m=>Number(m.dl_rate)||0,
                ul_rate:m=>Number(m.ul_rate)||0,
                downloading:m=>Number(m.downloading)||0,
                disk:m=>Number((m.disk||{}).percent)||0};
  const getm = mkey[ms.col] || mkey.name;
  const cards = list.sort((a,b)=>cmpVal(getm(a), getm(b), ms.dir)).map(m=>{
    const d = m.disk||{}; const has = d && Number(d.total)>0;
    const pct = Number(d.percent)||0; const pctTxt = (pct*100).toFixed(0);
    const cls = pct>=DISK_PAUSE?'disk-danger':(pct>=DISK_HIGH?'disk-high':'');
    const barCls = pct>=DISK_PAUSE?'danger':(pct>=DISK_HIGH?'high':'');
    const tag = pct>=DISK_PAUSE
        ? '<span class="mtag danger">磁盘占用高·已暂停</span>'
        : (pct>=DISK_HIGH?'<span class="mtag high">磁盘占用高</span>':'');
    const usage = has? (fmtSize(d.used)+' / '+fmtSize(d.total)) : '磁盘信息未上报';
    const cpu = (m.cpu!=null)? Number(m.cpu) : null;
    const mem = (m.mem!=null)? Number(m.mem) : null;
    const barc = (p)=> p>=90?'danger':(p>=80?'high':'');
    const statLine = (cpu!=null||mem!=null) ? `
      <div class="mstat">
        <div class="ms"><span class="mslab">CPU ${cpu!=null?cpu.toFixed(0)+'%':'-'}</span>
          <div class="mbar ${cpu!=null?barc(cpu):''}"><i style="width:${cpu!=null?cpu:0}%"></i></div></div>
        <div class="ms"><span class="mslab">内存 ${mem!=null?mem.toFixed(0)+'%':'-'}</span>
          <div class="mbar ${mem!=null?barc(mem):''}"><i style="width:${mem!=null?mem:0}%"></i></div></div>
      </div>` : '';
    const selCls = (filter.machine && filter.machine===m.name) ? ' sel' : '';
    const dlN = (m.downloading!=null)? m.downloading : 0;
    const dlR = Number(m.dl_rate)||0, ulR = Number(m.ul_rate)||0;
    return `<div class="mcard ${cls}${selCls}" onclick="selectMachine('${esc(m.name)}')"
        title="点击只看该机器的任务">
      <div class="mname">${esc(m.name)}${tag}</div>
      <div class="mbar ${barCls}"><i style="width:${has?pctTxt:0}%"></i></div>
      <div class="mmeta"><span>磁盘 ${has?pctTxt+'%':'-'}</span>
        <span>${esc(usage)}</span><span>下载中 ${dlN}</span></div>
      <div class="mmeta"><span>↓${fmtSpeed(dlR)}</span><span>↑${fmtSpeed(ulR)}</span></div>
      ${statLine}
    </div>`;
  }).join('');
  const hint = filter.machine
    ? `<span class="mfilter-tag" onclick="selectMachine('')">当前只看：${esc(filter.machine)} ✕</span>`
    : '';
  const msBtn = (label, col) => {
    const active = (machineSort||{}).col===col;
    return `<button class="msort${active?' on':''}" onclick="sortMachines('${col}')">`
         + `${label}${arrowFor(active, machineSort&&machineSort.dir)}</button>`;
  };
  const sortBar = `<span class="msortbar">排序：${msBtn('名称','name')}`
    + `${msBtn('下载速度','dl_rate')}${msBtn('上传速度','ul_rate')}`
    + `${msBtn('下载中任务','downloading')}${msBtn('磁盘','disk')}</span>`;
  $('machines').innerHTML = '<p class="mtitle">机器磁盘占用（点击卡片筛选该机任务）'
    + hint + sortBar + '</p>'
    + '<div class="machines">' + cards + '</div>';
}

function renderAlerts(machines){
  const items = [];
  Object.values(machines).forEach(m=>{
    const d = m.disk||{}; if(!(d && Number(d.total)>0)) return;
    const pct = Number(d.percent)||0;
    if(pct < DISK_HIGH) return;
    items.push({name:m.name, pct:pct, used:d.used, total:d.total, danger: pct>=DISK_PAUSE});
  });
  if(!items.length){ $('alerts').innerHTML=''; return; }
  const rows = items.sort((a,b)=>b.pct-a.pct).map(it=>{
    const cls = it.danger?'danger':'high';
    const lvl = it.danger?'已暂停':'告警';
    const p = (it.pct*100).toFixed(0);
    const msg = it.danger
      ? `磁盘占用 ${p}%，已超过 ${Math.round(DISK_PAUSE*100)}%，已暂停该机任务`
      : `磁盘占用 ${p}%，已超过 ${Math.round(DISK_HIGH*100)}%，占用偏高`;
    return `<div class="alert ${cls}"><span class="lvl">${lvl}</span>
      <span class="amac">${esc(it.name)}</span>
      <span class="amsg">${msg}（${fmtSize(it.used)} / ${fmtSize(it.total)}）</span></div>`;
  }).join('');
  $('alerts').innerHTML = '<p class="mtitle">告警消息</p>'
    + '<div class="alerts">' + rows + '</div>';
}

let toastTimer=null;
function toast(msg){
  const t=$('toast'); t.textContent=msg; t.classList.add('show');
  clearTimeout(toastTimer); toastTimer=setTimeout(()=>t.classList.remove('show'),1600);
}

async function ctl(tk, action){
  try{
    await fetch('/api/control',{method:'POST',
      headers:{'Content-Type':'application/json'},
      body:JSON.stringify({task_key:tk,action:action})});
    const label={pause:'已暂停',start:'已开始',delete:'已删除',
      force_ingest:'已触发强制入库'}[action]||action;
    toast(label); poll();
  }catch(e){ toast('操作失败'); }
}

// ---------------- 初始化 ----------------
resetGroups();
buildSkeleton();
initEvents();
['f-rid','f-start','f-end'].forEach(id=>$(id).addEventListener('change', onFilterChange));
// 搜索框：输入防抖 400ms 触发筛选；回车立即触发
let _qTimer=null;
$('f-q').addEventListener('input', ()=>{ clearTimeout(_qTimer);
  _qTimer=setTimeout(onFilterChange, 400); });
$('f-q').addEventListener('keydown', (e)=>{ if(e.key==='Enter'){
  clearTimeout(_qTimer); onFilterChange(); } });
poll();
setInterval(poll, POLL_MS);
</script>
</body></html>"""


class _Handler(BaseHTTPRequestHandler):
    def log_message(self, format, *args):  # 静默默认访问日志
        pass

    def _send_json(self, code: int, obj: dict):
        data = json.dumps(obj, ensure_ascii=False).encode("utf-8")
        self.send_response(code)
        self.send_header("Content-Type", "application/json; charset=utf-8")
        self.send_header("Content-Length", str(len(data)))
        self.end_headers()
        self.wfile.write(data)

    def _send_html(self, html: str):
        data = html.encode("utf-8")
        self.send_response(200)
        self.send_header("Content-Type", "text/html; charset=utf-8")
        self.send_header("Content-Length", str(len(data)))
        self.end_headers()
        self.wfile.write(data)

    def _send_csv(self, filename: str, text: str):
        data = text.encode("utf-8")
        self.send_response(200)
        self.send_header("Content-Type", "text/csv; charset=utf-8")
        self.send_header("Content-Disposition", f'attachment; filename="{filename}"')
        self.send_header("Content-Length", str(len(data)))
        self.end_headers()
        self.wfile.write(data)

    def _send_jsonl(self, filename: str, text: str):
        data = text.encode("utf-8")
        self.send_response(200)
        self.send_header("Content-Type", "application/x-ndjson; charset=utf-8")
        self.send_header("Content-Disposition", f'attachment; filename="{filename}"')
        self.send_header("Content-Length", str(len(data)))
        self.end_headers()
        self.wfile.write(data)

    def _read_body(self) -> dict:
        length = int(self.headers.get("Content-Length", 0) or 0)
        if length <= 0:
            return {}
        try:
            return json.loads(self.rfile.read(length).decode("utf-8"))
        except Exception:
            return {}

    @staticmethod
    def _q1(qs: dict, key: str, default=None):
        v = qs.get(key, [default])
        return v[0] if v else default

    def _handle_export(self, qs: dict):
        """导出 JSONL：每行 = 原始种子（raw_seed 原文）+ 现有逻辑状态（挂在 _status 字段）。

        原始种子按其原本 JSON 结构展开为顶层字段（便于直接重投/复用）；解析失败则原文
        放入 _raw_seed。当前逻辑判定的状态（分类/状态文案/进度/机器/速度/错误等）统一挂
        在 _status 字段，不覆盖原种子字段。
        """
        category = str(self._q1(qs, "category", "all") or "all")
        rid = self._q1(qs, "requirement_id", None)
        start = _to_float(self._q1(qs, "start", None))
        end = _to_float(self._q1(qs, "end", None))
        machine = self._q1(qs, "machine", None)
        q = self._q1(qs, "q", None)
        rows = STORE.export_seed_rows(category, rid, start, end,
                                      machine=machine, q=q)
        buf = io.StringIO()
        for raw, status in rows:
            seed_obj: Dict[str, Any] = {}
            try:
                parsed = json.loads(raw) if raw else {}
                if isinstance(parsed, dict):
                    seed_obj = parsed
                else:
                    seed_obj = {"_raw_seed": parsed}
            except Exception:
                seed_obj = {"_raw_seed": raw}
            seed_obj["_status"] = status
            buf.write(json.dumps(seed_obj, ensure_ascii=False) + "\n")
        fname = f"torrent_{category}_{int(time.time())}.jsonl"
        self._send_jsonl(fname, buf.getvalue())

    def do_GET(self):
        parsed = urlparse(self.path)
        qs = parse_qs(parsed.query)
        path = parsed.path
        if path in ("/", "/index.html"):
            self._send_html(_render_dashboard())
        elif path == "/api/summary":
            self._send_json(200, STORE.summary(
                self._q1(qs, "requirement_id", None),
                _to_float(self._q1(qs, "start", None)),
                _to_float(self._q1(qs, "end", None)),
                machine=self._q1(qs, "machine", None),
                q=self._q1(qs, "q", None)))
        elif path == "/api/list":
            try:
                offset = int(self._q1(qs, "offset", 0) or 0)
            except (TypeError, ValueError):
                offset = 0
            try:
                limit = int(self._q1(qs, "limit", 20) or 20)
            except (TypeError, ValueError):
                limit = 20
            limit = max(0, min(limit, 1000000))
            try:
                sort_dir = int(self._q1(qs, "dir", -1) or -1)
            except (TypeError, ValueError):
                sort_dir = -1
            self._send_json(200, STORE.list_page(
                str(self._q1(qs, "category", "all") or "all"),
                self._q1(qs, "requirement_id", None),
                _to_float(self._q1(qs, "start", None)),
                _to_float(self._q1(qs, "end", None)),
                offset, limit,
                machine=self._q1(qs, "machine", None),
                q=self._q1(qs, "q", None),
                sort=self._q1(qs, "sort", None),
                sort_dir=sort_dir))
        elif path == "/api/task":
            task = STORE.get_task((self._q1(qs, "task_key", "") or "").strip())
            if task is None:
                self._send_json(404, {"error": "not found"})
            else:
                self._send_json(200, {"task": task})
        elif path == "/api/events":
            try:
                offset = int(self._q1(qs, "offset", 0) or 0)
            except (TypeError, ValueError):
                offset = 0
            try:
                limit = int(self._q1(qs, "limit", 30) or 30)
            except (TypeError, ValueError):
                limit = 30
            limit = max(0, min(limit, 1000))
            self._send_json(200, STORE.list_events(
                self._q1(qs, "requirement_id", None),
                _to_float(self._q1(qs, "start", None)),
                _to_float(self._q1(qs, "end", None)),
                str(self._q1(qs, "level", "") or ""),
                offset, limit,
                machine=self._q1(qs, "machine", None),
                q=self._q1(qs, "q", None)))
        elif path == "/api/export":
            self._handle_export(qs)
        elif path == "/api/status":
            self._send_json(200, {"status": STORE.all_status()})
        elif path == "/api/dispatch":
            info_hash = (self._q1(qs, "info_hash", "") or "").strip()
            try:
                rid_q = int(self._q1(qs, "requirement_id", 0) or 0)
            except (TypeError, ValueError):
                rid_q = 0
            caller_machine = (self._q1(qs, "machine", "") or "").strip()
            self._send_json(200, STORE.dispatch_decision(
                info_hash, rid_q, caller_machine=caller_machine))
        elif path == "/api/pending_retries":
            machine = (self._q1(qs, "machine", "") or "").strip()
            self._send_json(200, {"seeds": STORE.claim_retries(machine)})
        else:
            self._send_json(404, {"error": "not found"})

    def do_POST(self):
        parsed = urlparse(self.path)
        body = self._read_body()
        if parsed.path == "/api/report":
            desired = STORE.report(body)
            self._send_json(200, {"desired": desired})
        elif parsed.path == "/api/control":
            ok = STORE.set_control(str(body.get("task_key") or ""),
                                   str(body.get("action") or ""))
            self._send_json(200 if ok else 400, {"ok": ok})
        elif parsed.path == "/api/drop_queued":
            ok = STORE.drop_task(str(body.get("task_key") or ""))
            self._send_json(200, {"ok": ok})
        elif parsed.path == "/api/retry":
            keys = body.get("task_keys")
            if not keys:
                one = str(body.get("task_key") or "")
                keys = [one] if one else []
            keys = [str(k) for k in keys if k]
            res = STORE.retry(keys)
            self._send_json(200 if res.get("ok") else 400, res)
        else:
            self._send_json(404, {"error": "not found"})


def serve(host: str = "0.0.0.0", port: int = PROCESS_SERVER_PORT,
          console_display: bool = True):
    """启动 HTTP 监控/控制服务（阻塞）。"""
    if console_display:
        threading.Thread(target=_console_display_loop, daemon=True).start()
    threading.Thread(target=_retry_sweeper_loop, daemon=True).start()
    server = ThreadingHTTPServer((host, port), _Handler)
    logger.info(f"[process] 监控服务启动: http://{host}:{port} "
                f"（对外 {PROCESS_SERVER_URL}），DB={DB_PATH}")
    try:
        server.serve_forever()
    except KeyboardInterrupt:
        logger.info("[process] 收到中断，关闭服务")
    finally:
        server.shutdown()


def _retry_sweeper_loop():
    """周期扫描本机重试队列，把原机长时间未认领的种子改投远端 job_id（防止永久积压）。"""
    try:
        interval = float(os.environ.get("TORRENT_RETRY_SWEEP_INTERVAL", "120") or 120)
    except ValueError:
        interval = 120.0
    try:
        max_age = float(os.environ.get("TORRENT_RETRY_STALE_SECONDS", "300") or 300)
    except ValueError:
        max_age = 300.0
    logger.info(f"[process] 重试队列滞留扫描已启动（每 {int(interval)}s 扫描，"
                f"滞留超 {int(max_age)}s 改投远端）")
    while True:
        time.sleep(interval)
        try:
            STORE.reroute_stale_retries(max_age)
        except Exception as e:
            logger.warning(f"[process] 重试队列滞留扫描异常: {e}")


def _console_display_loop(interval: float = 10.0):
    while True:
        time.sleep(interval)
        try:
            status = STORE.all_status()
            if not status:
                continue
            total_rate = sum(float(s.get("download_rate") or 0)
                             for s in status.values()
                             if s.get("state") == ST_DOWNLOADING and not s.get("offline"))
            lines = [f"===== 监控总览 任务={len(status)} 总速度={_fmt_speed(total_rate)} "
                     f"{time.strftime('%H:%M:%S')} ====="]
            for _key, s in sorted(status.items()):
                files = s.get("files") or []
                completed = sum(1 for f in files if f.get("status") == ST_COMPLETED)
                st_txt = "离线" if s.get("offline") else s.get("state_label") or s.get("state", "")
                ih = str(s.get("info_hash") or "")
                lines.append(
                    f"  {ih[:12]} rid={s.get('requirement_id',0)} "
                    f"{(s.get('name') or '')[:24]} "
                    f"{float(s.get('progress') or 0)*100:.1f}% "
                    f"{_fmt_speed(s.get('download_rate') or 0)} "
                    f"seeds={s.get('num_seeds',0)} peers={s.get('num_peers',0)} "
                    f"files={completed}/{s.get('num_files', len(files))} "
                    f"{st_txt}/{s.get('desired','')}")
            logger.info("\n".join(lines))
        except Exception as e:
            logger.warning(f"[process] 控制台展示异常: {e}")


# ============================================================================
# 客户端（downloader 侧使用）
# ============================================================================

class ProcessReporterClient:
    """downloader 用来向 process 服务上报状态并获取控制期望态。"""

    def __init__(self, server_url: str = PROCESS_SERVER_URL, timeout: float = 8.0):
        self.server_url = server_url.rstrip("/")
        self.timeout = timeout
        self._session = requests.Session()
        self._warned = False

    def report(self, snapshot: dict) -> str:
        """上报快照，返回该任务的 desired 状态；服务不可达时返回 running。"""
        try:
            resp = self._session.post(
                self.server_url + "/api/report",
                json=snapshot, timeout=self.timeout)
            if resp.status_code == 200:
                return str(resp.json().get("desired") or DESIRED_RUNNING)
        except Exception as e:
            if not self._warned:
                logger.warning(f"[reporter] 无法连接监控服务 {self.server_url}: {e}"
                               f"（后续将静默重试）")
                self._warned = True
        return DESIRED_RUNNING

    def report_queued(self, rid, qid: str, name: str = "", machine: str = "",
                      total_size: int = 0, num_files: int = 0) -> None:
        """上报一个“排队等待下载”占位（尚未开始处理的种子，key=rid:q:qid）。"""
        snap = {
            "info_hash": f"q:{qid}", "requirement_id": rid, "name": name,
            "machine": machine or "", "state": ST_QUEUED,
            "progress": 0.0, "download_rate": 0, "upload_rate": 0,
            "num_files": int(num_files or 0), "completed_files": 0,
            "total_size": int(total_size or 0),
        }
        try:
            self._session.post(self.server_url + "/api/report",
                               json=snap, timeout=self.timeout)
        except Exception:
            pass

    def drop_queued(self, rid, qid: str) -> None:
        """种子真正开始下载后，删除其排队占位。"""
        try:
            self._session.post(self.server_url + "/api/drop_queued",
                               json={"task_key": _task_key(rid, f"q:{qid}")},
                               timeout=self.timeout)
        except Exception:
            pass

    def control(self, info_hash: str, action: str, rid=0) -> bool:
        try:
            resp = self._session.post(
                self.server_url + "/api/control",
                json={"task_key": _task_key(rid, info_hash), "action": action},
                timeout=self.timeout)
            return resp.status_code == 200
        except Exception:
            return False

    def dispatch_decision(self, info_hash: str, rid=0, machine="") -> str:
        """询问该 (rid, info_hash) 是否应处理：返回 'proceed' / 'ignore' / 'completed'。

        machine 为本机 ID：服务端据此判断“请求是否来自原下载机”——若原机有未下完的本地
        文件且请求来自别的机器，则交回原机续传并让本机 ignore（复用本地、不重复下载）。
        服务不可达时返回 'proceed'（不因监控缺失而阻塞下载）。
        """
        try:
            resp = self._session.get(
                self.server_url + "/api/dispatch",
                params={"info_hash": info_hash, "requirement_id": rid,
                        "machine": machine},
                timeout=self.timeout)
            if resp.status_code == 200:
                return str(resp.json().get("decision") or "proceed")
        except Exception:
            pass
        return "proceed"

    def claim_retries(self, machine: str) -> list:
        """取回属于本机的待重试原始种子（原机本地重跑用）；服务不可达返回空列表。"""
        try:
            resp = self._session.get(
                self.server_url + "/api/pending_retries",
                params={"machine": machine}, timeout=self.timeout)
            if resp.status_code == 200:
                return list(resp.json().get("seeds") or [])
        except Exception:
            pass
        return []


if __name__ == "__main__":
    serve()
