# 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"}
  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 删除

# 状态 -> 中文展示文案（看板用）
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: "因存储不足已下载入库后删除",
}

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

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

# 看板分类顺序（key, 中文, marker class）
GROUP_KEYS = ["downloading", "seeding", "paused", "stalled", "force_ingested",
              "canceled_no_space", "deleted_no_space", "deleted", "other"]


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) -> str:
    """把 (state, desired) 归到看板分类（与前端保持一致，供 list/summary/export 共用）。

    未识别的状态一律落入 "other"（其他特殊情况），保证任何种子都能在前端显示、不丢失。
    """
    if state == ST_SEEDING:
        return "seeding"
    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):
        return "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()

        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._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)")
        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 _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) -> dict:
        """按 requirement_id + 时间范围（事件发生时间）+ 级别 分页查询历史事件。"""
        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))
        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 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]) -> 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
        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),
            "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) -> dict:
        counts = {k: 0 for k in GROUP_KEYS}
        total_rate = 0.0
        seeding_count = 0
        seed_up_rate = 0.0
        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)                       # 候选下拉：全量（不受筛选影响）
                if not self._passes(snap, rid, start, end):
                    continue
                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)
                counts[cat] = counts.get(cat, 0) + 1
                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)
                mn = snap.get("machine") or "(未知)"
                m = machines.setdefault(mn, {"name": mn, "tasks": 0, "disk": None})
                m["tasks"] += 1
                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
        return {
            "counts": counts, "total_rate": total_rate,
            "seeding_count": seeding_count, "seed_up_rate": seed_up_rate,
            "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) -> 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):
                    continue
                desired = self._desired.get(key, DESIRED_RUNNING)
                cat = _category_of(snap.get("state", ""), desired)
                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) -> List[dict]:
        return self.list_page(category, rid, start, end, offset=0, limit=0)["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}
.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}
.btns{display:flex;gap:6px}
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}
.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}
.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}
.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="添加结束时间">
    <button onclick="clearFilters()">清除</button>
  </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">-</div></div>
    <div class="stat"><div class="k">做种中</div><div class="v" id="s-seed">-</div></div>
    <div class="stat"><div class="k">已暂停</div><div class="v" id="s-pause">-</div></div>
    <div class="stat"><div class="k">停滞放弃</div><div class="v" id="s-stall">-</div></div>
    <div class="stat"><div class="k">强制入库</div><div class="v" id="s-ingest">-</div></div>
    <div class="stat"><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'],
  ['seeding','做种中','mk-ingest'],
  ['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'],
];
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:''};
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);
  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'};
  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 / 计数） ----------------
function buildSkeleton(){
  const html = GROUPS.map(([key,label,mk])=>{
    const resume = (key==='paused'||key==='stalled')
      ? `<button class="resume-all" onclick="event.stopPropagation();resumeAll('${key}')">全部继续</button>` : '';
    return `<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>${resume}
        <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 rate = (r.state==='downloading' && !r.offline)? (Number(r.download_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 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">${rate>0?fmtSpeed(rate):'-'}</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>`}
      <button class="danger" onclick="ctl('${esc(tk)}','delete')">删除</button>
    </div></td>
  </tr>`;
  const detail = open ? `<tr class="detail"><td colspan="8">${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="8">无</td></tr>`;
  } else {
    body = g.rows.map(r=>taskRow(r)).join('');
    if(!g.done)
      body += `<tr class="more"><td colspan="8">`
            + `<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="8">已全部加载 ${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>
  </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').textContent = groupTotals.downloading||0;
  $('s-seed').textContent = (d.seeding_count||0)
    + ((d.seed_up_rate>0)? (' · ↑'+fmtSpeed(d.seed_up_rate)) : '');
  $('s-pause').textContent = groupTotals.paused||0;
  $('s-stall').textContent = groupTotals.stalled||0;
  $('s-ingest').textContent = groupTotals.force_ingested||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();
}

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='';
  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) : '';
}
function onFilterChange(){
  readFilters(); resetGroups(); buildSkeleton();
  evState = {rows:[], offset:0, total:0, done:false, loading:false, _sig:null};
  initEvents();
  poll();
}

// ---------------- 消息列表（历史错误 / 告警，按 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)) : '磁盘信息未上报';
    return `<div class="mcard ${cls}">
      <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>
    </div>`;
  }).join('');
  $('machines').innerHTML = '<p class="mtitle">机器磁盘占用</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));
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))
        rows = STORE.export_rows(category, rid, start, end)
        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))))
        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))
        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))
        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))
        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})
        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"


if __name__ == "__main__":
    serve()
