# 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 csv
import datetime
import io
import json
import os
import sqlite3
import threading
import time
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
from typing import 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_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)，本地已清

# 下载器自身 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_FORCE_INGESTED: "已强制入库",
    ST_SEEDING: "做种中",
    ST_CANCELED_NO_SPACE: "因存储不足取消下载",
    ST_DELETED_NO_SPACE: "因存储不足已下载入库后删除",
    ST_REQUEUED: "重投队列",
    ST_INGESTED_DONE: "已全部入库完成",
}

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

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

# 看板分类顺序（key）。做种中放最后（页面最底部）；重投队列在其上。
GROUP_KEYS = ["downloading", "paused", "stalled", "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_SEEDING:
        return "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 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），用于“搜索某种子后只看它的报错消息”。
        """
        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 info_hash=?)")
                args.extend([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)
            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
            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 dispatch_decision(self, info_hash: str, rid=0) -> dict:
        """重复投递同一 (requirement_id, info_hash) 时的处理决策（下载端开始前询问）。

        去重按 (rid, info_hash)：同一需求下重复投递才 ignore/completed；不同需求复用
        同一种子各自独立处理（本机已有完整数据则 downloader 走本地直接入库快路径）。
        规则：
          - 本 rid 已完成                          -> completed（跳过）
          - 本 rid 在线下载中/元数据/上传/等待      -> ignore（存活实例在跑）
          - 本 rid 在线已暂停（存活实例）           -> ignore + desired=running（让其继续）
          - 停滞放弃 / 错误                        -> 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()
        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
                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"}
                # 到这里都会 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 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, "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
                # 计数 / 速度 / 流量：叠加机器过滤（选中机器时只统计该机）
                if machine not in (None, "", "all") and mn != machine:
                    continue
                m["tasks"] += 1
                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)
                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) -> 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))
            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 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)}
.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-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('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('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 = [
  ['downloading','正在下载','mk-dl'],
  ['paused','已暂停','mk-pause'],
  ['stalled','下载停滞放弃','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();   // 折叠的分类
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:''};
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;               // 消息列表默认收起，点按钮展开

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;}

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',force_ingested:'b-ingested',
    canceled_no_space:'b-error',deleted_no_space:'b-offline',requeued:'b-pending',
    ingested_done:'b-completed'};
  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','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?'':`<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 flimit = fileLimits[tk] || PAGE;
  const shownFiles = files.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||'';
    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 style="min-width:120px"><div class="bar sm ${done?'done':''}">
        <i style="width:${fpct}%"></i></div></td>
      <td>${fpct}%</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>`
      : '';
  return `<div class="det">
    <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>
        <th>相对路径</th><th>模态</th><th>大小</th><th>进度</th><th></th><th>状态</th></tr></thead>
        <tbody>${fileRows}${moreFiles}</tbody></table></div>
    </div>
  </div>`;
}

// 分类内容签名：无变化则跳过重绘（静默刷新防卡顿的核心）
function groupSig(key){
  const g = groupState[key];
  let s = key+'|'+g.total+'|'+g.done+'|'+foldedGroups.has(key);
  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);
  }
  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 {
    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>`;
  }
  cont.innerHTML = `<table><thead><tr>
    <th>info_hash</th><th>机器</th><th>状态</th><th>进度</th>
    <th>速度(↓/↑)</th><th>总大小</th><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({category:key, offset:off, limit:PAGE}),{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({category:key, offset:0, limit:cnt}),{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-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-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?'+fq({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 =
      '<button id="ev-toggle" class="ev-toggle" onclick="toggleEvents()">'
    + '<span class="gcaret">&#9656;</span>消息列表（历史错误 / 告警）'
    + '<span class="ev-cnt" id="ev-cnt"></span></button>'
    + '<div class="evlist" id="ev-list" style="display:none" '
    + 'onscroll="onEventsScroll(event)"></div>';
  renderEvents(true);
}

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?'+fq({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?'+fq({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） ----------------
function renderMachines(machines){
  const list = Object.values(machines);
  if(!list.length){ $('machines').innerHTML=''; return; }
  const cards = list.sort((a,b)=>a.name.localeCompare(b.name)).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' : '';
    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>任务 ${m.tasks}</span></div>
      ${statLine}
    </div>`;
  }).join('');
  const hint = filter.machine
    ? `<span class="mfilter-tag" onclick="selectMachine('')">当前只看：${esc(filter.machine)} ✕</span>`
    : '';
  $('machines').innerHTML = '<p class="mtitle">机器磁盘占用（点击卡片筛选该机任务）'
    + hint + '</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 _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):
        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_rows(category, rid, start, end, machine=machine, q=q)
        buf = io.StringIO()
        buf.write("\ufeff")   # BOM，便于 Excel 正确识别中文
        w = csv.writer(buf)
        w.writerow(["info_hash", "requirement_id", "name", "machine", "state",
                    "progress", "download_rate_Bps", "upload_rate_Bps",
                    "num_seeds", "num_peers", "completed_files", "num_files",
                    "total_size", "added_time", "updated_time"])

        def _ts(t):
            try:
                return datetime.datetime.fromtimestamp(float(t)).strftime(
                    "%Y-%m-%d %H:%M:%S") if t else ""
            except Exception:
                return ""

        for r in rows:
            w.writerow([
                r["info_hash"], r["requirement_id"], r["name"], r["machine"],
                r["state_label"], f'{r["progress"]*100:.1f}%',
                int(r["download_rate"]), int(r["upload_rate"]),
                r["num_seeds"], r["num_peers"], r["completed_files"],
                r["num_files"], r["total_size"], _ts(r["added_ts"]),
                _ts(r["updated_ts"])])
        fname = f"torrent_{category}_{int(time.time())}.csv"
        self._send_csv(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))
            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)))
        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
            self._send_json(200, STORE.dispatch_decision(info_hash, rid_q))
        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/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()
    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 _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 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) -> str:
        """询问该 (rid, info_hash) 是否应处理：返回 'proceed' / 'ignore' / 'completed'。

        服务不可达时返回 'proceed'（不因监控缺失而阻塞下载）。响应里还带 mode
        （continue/restart/skip/new）作为附加信息，此处仅回传 decision 字符串保持
        与既有 downloader 调用兼容。
        """
        try:
            resp = self._session.get(
                self.server_url + "/api/dispatch",
                params={"info_hash": info_hash, "requirement_id": rid},
                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()
