import html
import json
import os
import re
import shutil
import subprocess
import threading
import time
from datetime import datetime
from pathlib import Path
from server.services import playlist as playlists
from server.services.covers import covers
from server.services.device import DeviceProfile
from server.services.ffmpeg_tool import convert_to_mp4, ensure_ffmpeg
from server.services.http import http
from server.services.live import live
from server.services.rapt import API_BASE
from server.services.sessions import sessions
ROOT = Path(__file__).resolve().parents[2]
STATE_DIR = ROOT / "cache" / "downloads"
STATE_PATH = STATE_DIR / "jobs.json"
SETTINGS_PATH = STATE_DIR / "settings.json"
DEFAULT_WORKERS = 6
MIN_WORKERS = 1
ACTIVE = {"queued", "downloading", "paused", "failed"}
FOLDER_NOTE = """这部剧的文件说明
封面
和本说明放在一起的封面图片,是这部剧的封面。文件名是封面.jpg 或封面.png,以实际图片格式为准。
简介.txt
这部剧的编号、名称和简介。
MP4
每一集由多个分片合成为一条完整的单集视频,放在这个目录里。一集对应一个文件,可以直接用播放器打开。
source
这里保存每一集下载下来的分片和播放列表。MP4 目录里的完整单集,就是由这些分片合成的。平时观看用 MP4 目录即可。
"""
def _plain_text(value: str) -> str:
text = html.unescape(str(value or ""))
text = re.sub(r"(?i)
", "\n", text)
text = re.sub(r"(?i)
", "\n", text)
text = re.sub(r"<[^>]+>", "", text)
text = text.replace("\r\n", "\n").replace("\r", "\n")
lines = [re.sub(r"[ \t]+", " ", line).strip() for line in text.split("\n")]
return "\n".join(line for line in lines if line).strip()
def _now() -> str:
return datetime.now().strftime("%Y-%m-%d %H:%M")
def _default_directory() -> str:
return str((Path.cwd() / "download").resolve())
class DownloadManager:
def __init__(self) -> None:
self._lock = threading.Lock()
self._jobs: dict[str, dict] = {}
self._directory = _default_directory()
self._worker_limit = DEFAULT_WORKERS
self._alive_workers = 0
self._wake = threading.Event()
self._started = False
self._speeds: dict[str, list[tuple[float, int]]] = {}
self._load()
def start(self) -> None:
with self._lock:
first = not self._started
self._started = True
if first:
for job in self._jobs.values():
if job["status"] in ("downloading", "failed"):
job["status"] = "queued"
job["error"] = ""
job["command"] = ""
for ep in job.get("episodes") or []:
if ep.get("phase") in ("failed", "downloading", "converting"):
ep["phase"] = "pending"
self._save_locked()
self._spawn_locked()
if first:
threading.Thread(target=self._prepare_ffmpeg, name="ffmpeg-setup", daemon=True).start()
def _spawn_locked(self) -> None:
while self._alive_workers < self._worker_limit:
self._alive_workers += 1
threading.Thread(target=self._worker, name="download-worker", daemon=True).start()
def _prepare_ffmpeg(self) -> None:
try:
ensure_ffmpeg()
except Exception:
return
def snapshot(self) -> dict:
with self._lock:
active = []
done = []
for job in self._jobs.values():
view = self._view_locked(job)
if job["status"] in ACTIVE:
active.append(view)
elif job["status"] == "completed":
done.append(view)
done.sort(key=lambda item: item.get("completed_at") or "", reverse=True)
return {
"directory": self._directory,
"workers": self._worker_limit,
"active": active,
"done": done,
"active_count": len(active),
}
def settings_view(self) -> dict:
with self._lock:
return {
"directory": self._directory,
"workers": self._worker_limit,
"default_directory": _default_directory(),
}
def set_workers(self, count: int) -> int:
count = int(count)
if count < MIN_WORKERS:
raise ValueError(f"同时下载数量至少为 {MIN_WORKERS}")
with self._lock:
self._worker_limit = count
self._save_settings_locked()
if self._started:
self._spawn_locked()
self._wake.set()
return count
def choose_directory(self) -> str:
return _pick_folder(self._directory)
def set_directory(self, path: str) -> str:
chosen = str(Path(path).expanduser().resolve())
Path(chosen).mkdir(parents=True, exist_ok=True)
with self._lock:
self._directory = chosen
self._save_settings_locked()
return chosen
def pick_directory(self) -> str:
chosen = _pick_folder(self._directory)
if not chosen:
return self._directory
return self.set_directory(chosen)
def enqueue(self, token: str, raw_device: dict | None, items: list[dict]) -> dict:
self.start()
device = sessions.bind(token, DeviceProfile.merge(raw_device)) if token else None
with self._lock:
directory = self._directory
Path(directory).mkdir(parents=True, exist_ok=True)
for item in items:
drama_id = str(item.get("id") or "").strip()
if not drama_id:
continue
self._enqueue_one(
drama_id,
str(item.get("title") or drama_id),
str(item.get("image") or ""),
token,
device,
directory,
)
self._wake.set()
return self.snapshot()
def pause_all(self) -> dict:
with self._lock:
changed = False
for job in self._jobs.values():
if job["status"] in ("queued", "downloading"):
job["status"] = "paused"
job["command"] = "pause"
changed = True
if changed:
self._save_locked()
return self.snapshot()
def resume_all(self) -> dict:
with self._lock:
changed = False
for job in self._jobs.values():
if job["status"] not in ("paused", "failed"):
continue
job["status"] = "queued"
job["command"] = ""
job["error"] = ""
for ep in job.get("episodes") or []:
if ep.get("phase") == "failed":
ep["phase"] = "pending"
changed = True
if changed:
self._save_locked()
if changed:
self._wake.set()
return self.snapshot()
def pause(self, drama_id: str) -> dict:
with self._lock:
job = self._jobs.get(drama_id)
if job and job["status"] in ("queued", "downloading"):
job["status"] = "paused"
job["command"] = "pause"
self._save_locked()
return self.snapshot()
def resume(self, drama_id: str) -> dict:
with self._lock:
job = self._jobs.get(drama_id)
if job and job["status"] in ("paused", "failed"):
job["status"] = "queued"
job["command"] = ""
job["error"] = ""
for ep in job.get("episodes") or []:
if ep.get("phase") == "failed":
ep["phase"] = "pending"
self._save_locked()
self._wake.set()
return self.snapshot()
def cancel(self, drama_id: str) -> dict:
with self._lock:
job = self._jobs.get(drama_id)
if job and job["status"] in ACTIVE:
job["command"] = "cancel"
if job["status"] != "downloading":
self._jobs.pop(drama_id, None)
self._save_locked()
self._wake.set()
return self.snapshot()
def remove(self, drama_id: str) -> dict:
with self._lock:
job = self._jobs.pop(drama_id, None)
self._save_locked()
folder = job.get("folder") if job else ""
if folder and Path(folder).is_dir():
shutil.rmtree(folder, ignore_errors=True)
return self.snapshot()
def open_folder(self, drama_id: str) -> None:
with self._lock:
job = self._jobs.get(drama_id)
folder = job.get("folder") if job else ""
if not folder or not Path(folder).is_dir():
raise RuntimeError("目录不存在")
os.startfile(folder) # noqa: S606
def _enqueue_one(self, drama_id: str, title: str, image: str, token: str, device, directory: str) -> None:
episodes = playlists.load_episodes(token, device, int(drama_id))
with self._lock:
job = self._jobs.get(drama_id)
if job is None:
folder = self._folder_locked(directory, title, drama_id)
job = {
"id": drama_id,
"title": title,
"image": image,
"status": "queued",
"command": "",
"folder": str(folder),
"created_at": _now(),
"completed_at": "",
"episodes": [],
"error": "",
"fraction": 0,
}
self._jobs[drama_id] = job
else:
if title:
job["title"] = title
if image:
job["image"] = image
if job["status"] == "completed":
job["status"] = "queued"
job["completed_at"] = ""
elif job["status"] == "paused":
pass
elif job["status"] not in ACTIVE:
job["status"] = "queued"
self._apply_episodes_locked(job, episodes)
if job["status"] == "paused":
job["command"] = "pause"
self._save_locked()
def _apply_episodes_locked(self, job: dict, episodes: list[dict]) -> None:
current = {int(ep["cid"]): ep for ep in job.get("episodes") or []}
merged = []
for ep in episodes:
old = current.get(int(ep["cid"]))
phase = old.get("phase") if old else "pending"
if phase not in ("done", "downloading", "converting"):
phase = "pending" if ep.get("src") else "skipped"
if phase == "done" and not self._mp4_exists(job, ep):
phase = "pending" if ep.get("src") else "skipped"
merged.append({
"cid": int(ep["cid"]),
"idx": int(ep.get("idx") or 0),
"title": ep.get("title") or "",
"src": ep.get("src") or "",
"phase": phase,
})
job["episodes"] = merged
if any(ep["phase"] == "pending" for ep in merged) and job["status"] == "completed":
job["status"] = "queued"
job["completed_at"] = ""
def _folder_locked(self, directory: str, title: str, drama_id: str) -> Path:
base = playlists.safe_name(title, f"drama-{drama_id}")
used = {Path(job["folder"]).name for job in self._jobs.values() if job.get("folder")}
name = base if base not in used else playlists.safe_name(f"{base} {drama_id}", base)
return Path(directory) / name
def _write_folder_notes(self, folder: Path, image: str, drama_id: str, title: str, token: str, device) -> None:
(folder / "说明.txt").write_text(FOLDER_NOTE, encoding="utf-8")
self._save_cover(folder, image)
self._save_intro(folder, drama_id, title, token, device)
def _save_intro(self, folder: Path, drama_id: str, title: str, token: str, device) -> None:
vid = drama_id
name = title.strip()
desc = ""
if token and device is not None:
try:
data = http.request_data(
"POST",
API_BASE + "/app/video/videoinfo",
token,
device,
{"vid": int(drama_id)},
)
video = data.get("video") if isinstance(data, dict) else None
if isinstance(video, dict):
if video.get("id") not in (None, ""):
vid = video.get("id")
fetched = str(video.get("title") or "").strip()
if fetched:
name = fetched
desc = _plain_text(video.get("desc") or "")
except Exception:
pass
(folder / "简介.txt").write_text(
f"ID: {vid}\n名称: {name}\n简介: {desc}\n",
encoding="utf-8",
)
def _save_cover(self, folder: Path, image: str) -> None:
digest = (image or "").rstrip("/").rsplit("/", 1)[-1].strip()
if not digest:
return
try:
data, media = covers.read(digest)
except Exception:
return
if not data:
return
kind = media.split(";", 1)[0].strip().lower()
ext = {
"image/jpeg": ".jpg",
"image/png": ".png",
"image/webp": ".webp",
"image/gif": ".gif",
}.get(kind, ".jpg")
target = folder / f"封面{ext}"
for old in folder.glob("封面.*"):
if old.is_file() and old.resolve() != target.resolve():
old.unlink()
target.write_bytes(data)
def _worker(self) -> None:
while True:
with self._lock:
if self._alive_workers > self._worker_limit:
self._alive_workers -= 1
return
drama_id = self._claim()
if not drama_id:
self._wake.wait(0.4)
self._wake.clear()
continue
try:
self._run(drama_id)
except Exception as exc:
with self._lock:
job = self._jobs.get(drama_id)
if job is None or job.get("status") == "paused" or job.get("command") in ("pause", "cancel"):
continue
job["status"] = "queued"
job["error"] = str(exc)
self._save_locked()
def _claim(self) -> str:
with self._lock:
for job in self._jobs.values():
if job["status"] == "queued":
job["status"] = "downloading"
job["command"] = ""
job["error"] = ""
self._save_locked()
return job["id"]
return ""
def _run(self, drama_id: str) -> None:
with self._lock:
job = self._jobs.get(drama_id)
if job is None:
return
token = sessions.active_token()
device = sessions.device_for(token) if token else None
folder = Path(job["folder"])
image = str(job.get("image") or "")
title = str(job.get("title") or "")
episodes = playlists.load_episodes(token, device, int(drama_id))
with self._lock:
job = self._jobs.get(drama_id)
if job is None:
return
self._apply_episodes_locked(job, episodes)
rows = [dict(ep) for ep in job["episodes"]]
self._save_locked()
folder.mkdir(parents=True, exist_ok=True)
self._write_folder_notes(folder, image, drama_id, title, token, device)
self._prepare_playlists(drama_id, folder, rows)
for row in rows:
if self._stopped(drama_id):
return
try:
self._download_episode(drama_id, folder, row)
except RuntimeError as exc:
if self._stopped(drama_id):
return
self._set_phase(drama_id, row["cid"], "failed")
with self._lock:
job = self._jobs.get(drama_id)
if job is not None:
job["error"] = str(exc)
self._save_locked()
with self._lock:
job = self._jobs.get(drama_id)
if job is None or job.get("command") == "cancel":
self._jobs.pop(drama_id, None)
self._save_locked()
return
pending = any(ep["phase"] == "pending" for ep in job["episodes"] if ep.get("src"))
failed = any(ep["phase"] == "failed" for ep in job["episodes"])
if job["status"] == "paused":
pass
elif pending:
job["status"] = "queued"
elif failed:
job["status"] = "failed"
else:
job["status"] = "completed"
job["completed_at"] = _now()
job["fraction"] = 1
job["error"] = ""
self._save_locked()
def _stopped(self, drama_id: str) -> bool:
with self._lock:
job = self._jobs.get(drama_id)
if job is None:
return True
if job.get("command") == "cancel":
self._jobs.pop(drama_id, None)
self._save_locked()
return True
if job.get("command") == "pause" or job["status"] == "paused":
job["status"] = "paused"
job["command"] = "pause"
self._save_locked()
return True
return False
def _prepare_playlists(self, drama_id: str, folder: Path, rows: list[dict]) -> None:
self._set_stage(drama_id, "统计清单")
try:
for row in rows:
if not row.get("src"):
continue
if self._stopped(drama_id):
raise RuntimeError("已停止")
try:
self._write_playlist(drama_id, folder, row)
except RuntimeError as exc:
if self._stopped(drama_id):
raise
self._set_phase(drama_id, row["cid"], "failed")
with self._lock:
job = self._jobs.get(drama_id)
if job is not None:
job["error"] = str(exc)
self._save_locked()
finally:
self._set_stage(drama_id, "")
def _episode_paths(self, folder: Path, row: dict) -> tuple[Path, Path]:
episode_name = playlists.safe_name(row.get("title") or "", f"EP.{row.get('idx') or row['cid']}")
source_dir = folder / "source" / episode_name
mp4_path = folder / "MP4" / f"{episode_name}.mp4"
return source_dir, mp4_path
def _write_playlist(self, drama_id: str, folder: Path, row: dict) -> None:
source_dir, _mp4_path = self._episode_paths(folder, row)
master_text, media_text, items = playlists.expand_variant(row["src"])
source_dir.mkdir(parents=True, exist_ok=True)
(source_dir / "video.m3u8").write_text(media_text, encoding="utf-8", newline="\n")
(source_dir / "playlist.m3u8").write_text(
_localize(master_text, ["video.m3u8"]),
encoding="utf-8",
newline="\n",
)
(source_dir / "segments.json").write_text(json.dumps(items, ensure_ascii=False), encoding="utf-8")
done = sum(1 for item in items if (source_dir / item["name"]).is_file() and (source_dir / item["name"]).stat().st_size > 0)
self._set_segments(drama_id, row["cid"], len(items), done)
def _set_segments(self, drama_id: str, cid: int, count: int, done: int) -> None:
with self._lock:
job = self._jobs.get(drama_id)
if job is None:
return
for ep in job["episodes"]:
if int(ep["cid"]) == int(cid):
ep["segment_count"] = count
ep["segments_done"] = done
self._save_locked()
def _set_stage(self, drama_id: str, stage: str) -> None:
with self._lock:
job = self._jobs.get(drama_id)
if job is not None:
job["stage"] = stage
live.touch()
def _download_episode(self, drama_id: str, folder: Path, row: dict) -> None:
if not row.get("src"):
self._set_phase(drama_id, row["cid"], "skipped")
return
source_dir, mp4_path = self._episode_paths(folder, row)
manifest = source_dir / "segments.json"
if mp4_path.is_file() and mp4_path.stat().st_size > 0:
if manifest.is_file():
items = json.loads(manifest.read_text(encoding="utf-8"))
self._set_segments(drama_id, row["cid"], len(items), len(items))
self._set_phase(drama_id, row["cid"], "done")
return
self._set_phase(drama_id, row["cid"], "downloading")
self._write_playlist(drama_id, folder, row)
items = json.loads(manifest.read_text(encoding="utf-8"))
total = len(items)
for index, item in enumerate(items, start=1):
if self._stopped(drama_id):
raise RuntimeError("已停止")
target = source_dir / item["name"]
if not target.is_file() or target.stat().st_size <= 0:
http.download_file(
item["url"],
target,
on_chunk=lambda size, job_id=drama_id: self._note_bytes(job_id, size),
stop=lambda job_id=drama_id: self._stopped(job_id),
)
done = sum(1 for part in items[:index] if (source_dir / part["name"]).is_file() and (source_dir / part["name"]).stat().st_size > 0)
earlier = sum(1 for part in items[index:] if (source_dir / part["name"]).is_file() and (source_dir / part["name"]).stat().st_size > 0)
self._set_segments(drama_id, row["cid"], total, done + earlier)
if self._stopped(drama_id):
raise RuntimeError("已停止")
self._set_phase(drama_id, row["cid"], "converting")
convert_to_mp4(source_dir, mp4_path)
self._set_phase(drama_id, row["cid"], "done")
self._set_segments(drama_id, row["cid"], total, total)
def _set_phase(self, drama_id: str, cid: int, phase: str) -> None:
with self._lock:
job = self._jobs.get(drama_id)
if job is None:
return
for ep in job["episodes"]:
if int(ep["cid"]) == int(cid):
ep["phase"] = phase
self._save_locked()
def _set_fraction(self, drama_id: str, fraction: float) -> None:
with self._lock:
job = self._jobs.get(drama_id)
if job is not None:
job["fraction"] = fraction
def _note_bytes(self, drama_id: str, size: int) -> None:
now = time.time()
window = self._speeds.setdefault(drama_id, [])
window.append((now, size))
cutoff = now - 2
self._speeds[drama_id] = [(stamp, nbytes) for stamp, nbytes in window if stamp >= cutoff]
def _speed(self, drama_id: str) -> str:
window = self._speeds.get(drama_id) or []
if len(window) < 1:
return ""
span = max(0.4, window[-1][0] - window[0][0])
return _format_speed(sum(size for _stamp, size in window) / span)
def _mp4_exists(self, job: dict, ep: dict) -> bool:
folder = job.get("folder") or ""
if not folder:
return False
name = playlists.safe_name(ep.get("title") or "", f"EP.{ep.get('idx') or ep.get('cid')}")
path = Path(folder) / "MP4" / f"{name}.mp4"
return path.is_file() and path.stat().st_size > 0
def _view_locked(self, job: dict) -> dict:
episodes = job.get("episodes") or []
downloadable = [ep for ep in episodes if ep.get("src")]
total = len(downloadable)
done = sum(1 for ep in downloadable if ep.get("phase") == "done")
parts = sum(int(ep.get("segment_count") or 0) for ep in downloadable)
got = sum(int(ep.get("segments_done") or 0) for ep in downloadable)
progress = round(got / parts * 100) if parts else 0
return {
"id": job["id"],
"title": job["title"],
"image": job.get("image") or "",
"status": job["status"],
"episode_count": total,
"done_count": done,
"progress": progress,
"speed": self._speed(job["id"]) if job["status"] == "downloading" else "",
"stage": job.get("stage") or "",
"completed_at": job.get("completed_at") or "",
"error": job.get("error") or "",
}
def _load(self) -> None:
STATE_DIR.mkdir(parents=True, exist_ok=True)
if SETTINGS_PATH.is_file():
try:
data = json.loads(SETTINGS_PATH.read_text(encoding="utf-8"))
if isinstance(data, dict) and data.get("directory"):
self._directory = data["directory"]
workers = data.get("workers") if isinstance(data, dict) else None
if isinstance(workers, int) and workers >= MIN_WORKERS:
self._worker_limit = workers
except (OSError, json.JSONDecodeError):
pass
if not STATE_PATH.is_file():
return
try:
data = json.loads(STATE_PATH.read_text(encoding="utf-8"))
except (OSError, json.JSONDecodeError):
return
jobs = data.get("jobs") if isinstance(data, dict) else None
if isinstance(jobs, list):
for job in jobs:
if isinstance(job, dict) and job.get("id"):
self._jobs[str(job["id"])] = job
def _save_locked(self) -> None:
STATE_DIR.mkdir(parents=True, exist_ok=True)
stored = []
for job in self._jobs.values():
stored.append({key: value for key, value in job.items() if key != "fraction"})
temporary = STATE_PATH.with_suffix(".json.tmp")
temporary.write_text(json.dumps({"jobs": stored}, ensure_ascii=False), encoding="utf-8")
temporary.replace(STATE_PATH)
self._save_settings_locked()
live.touch()
def _save_settings_locked(self) -> None:
STATE_DIR.mkdir(parents=True, exist_ok=True)
SETTINGS_PATH.write_text(
json.dumps({"directory": self._directory, "workers": self._worker_limit}, ensure_ascii=False),
encoding="utf-8",
)
def _localize(text: str, names: list[str]) -> str:
lines = []
index = 0
for line in text.splitlines():
stripped = line.strip()
if stripped and not stripped.startswith("#"):
lines.append(names[index] if index < len(names) else stripped)
index += 1
else:
lines.append(line)
return "\n".join(lines) + "\n"
def _format_speed(value: float) -> str:
if value >= 1024 * 1024:
return f"{value / 1024 / 1024:.1f} MB/s"
if value >= 1024:
return f"{value / 1024:.0f} KB/s"
if value > 0:
return f"{value:.0f} B/s"
return ""
def _pick_folder(initial: str) -> str:
script = """
Add-Type -AssemblyName System.Windows.Forms
$owner = New-Object System.Windows.Forms.Form
$owner.TopMost = $true
$owner.ShowInTaskbar = $false
$owner.Opacity = 0
$owner.Show()
$owner.Activate()
$dialog = New-Object System.Windows.Forms.FolderBrowserDialog
$dialog.Description = '选择下载目录'
$dialog.ShowNewFolderButton = $true
$dialog.SelectedPath = $env:RAPTD_PICK_INITIAL
$result = $dialog.ShowDialog($owner)
$owner.Close()
if ($result -eq [System.Windows.Forms.DialogResult]::OK) {
[Console]::OutputEncoding = New-Object System.Text.UTF8Encoding $false
[Console]::Out.WriteLine($dialog.SelectedPath)
}
"""
env = os.environ.copy()
env["RAPTD_PICK_INITIAL"] = initial
completed = subprocess.run(
["powershell", "-NoProfile", "-STA", "-Command", script],
capture_output=True,
text=True,
encoding="utf-8",
env=env,
check=False,
)
return (completed.stdout or "").strip()
downloads = DownloadManager()