import json import os import shutil import subprocess import threading import time from datetime import datetime from pathlib import Path from urllib.parse import urlsplit from server.services import playlist as playlists 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.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 = 2 MIN_WORKERS = 1 MAX_WORKERS = 8 ACTIVE = {"queued", "downloading", "paused", "failed"} 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 or count > MAX_WORKERS: raise ValueError(f"同时下载数量需要在 {MIN_WORKERS} 到 {MAX_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 _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"]) 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._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, segments = playlists.expand_variant(row["src"]) items = [] for index, segment in enumerate(segments, start=1): name = Path(urlsplit(segment).path).name or f"part-{index}.ts" items.append({"name": name, "url": segment}) source_dir.mkdir(parents=True, exist_ok=True) (source_dir / "video.m3u8").write_text( _localize(media_text, [item["name"] for item in items]), 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") if not manifest.is_file(): 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 MIN_WORKERS <= workers <= MAX_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()