feat: 补齐桌面端下载播放,并将窗口改为夜间主题和应用图标

This commit is contained in:
Nixevol
2026-09-27 06:16:50 +08:00
parent f010426f94
commit 1b33898fff
81 changed files with 2529 additions and 232 deletions
+11
View File
@@ -10,8 +10,12 @@ from server.routes.auth import router as auth_router
from server.routes.catalog import router as catalog_router
from server.routes.covers import router as covers_router
from server.routes.device import router as device_router
from server.routes.downloads import router as downloads_router
from server.routes.live import router as live_router
from server.routes.play import router as play_router
from server.routes.proxy import router as proxy_router
from server.routes.session import router as session_router
from server.services.downloads import downloads
def create_app() -> FastAPI:
@@ -28,8 +32,15 @@ def create_app() -> FastAPI:
app.include_router(device_router)
app.include_router(catalog_router)
app.include_router(covers_router)
app.include_router(downloads_router)
app.include_router(play_router)
app.include_router(live_router)
app.include_router(proxy_router)
@app.on_event("startup")
def start_downloads() -> None:
downloads.start()
@app.get("/health")
def health():
return {"ok": True}
+1
View File
@@ -1,3 +1,4 @@
fastapi>=0.115
uvicorn>=0.32
websockets>=14
httpx[socks]>=0.28
+112
View File
@@ -0,0 +1,112 @@
from fastapi import APIRouter, HTTPException
from pydantic import BaseModel
from server.errors import call
from server.schemas import DeviceBody, device_payload
from server.services.downloads import downloads
router = APIRouter(prefix="/api/downloads")
class DownloadItem(BaseModel):
id: int
title: str = ""
image: str = ""
class DownloadStartBody(BaseModel):
token: str = ""
device: DeviceBody | None = None
items: list[DownloadItem] = []
class DirectoryBody(BaseModel):
path: str = ""
class DownloadSettingsBody(BaseModel):
directory: str = ""
workers: int | None = None
@router.get("")
def list_downloads():
return downloads.snapshot()
@router.post("")
def start_downloads(body: DownloadStartBody):
return downloads.enqueue(body.token, device_payload(body.device), [item.model_dump() for item in body.items])
@router.post("/directory")
def set_directory(body: DirectoryBody):
if not body.path.strip():
raise HTTPException(status_code=400, detail="请选择目录")
return {"directory": downloads.set_directory(body.path.strip())}
@router.post("/pick")
def pick_directory():
return {"directory": downloads.pick_directory()}
@router.post("/choose")
def choose_directory():
return {"directory": downloads.choose_directory()}
@router.get("/settings")
def get_download_settings():
return downloads.settings_view()
@router.post("/settings")
def update_download_settings(body: DownloadSettingsBody):
if body.directory.strip():
downloads.set_directory(body.directory.strip())
if body.workers is not None:
try:
downloads.set_workers(body.workers)
except ValueError as exc:
raise HTTPException(status_code=400, detail=str(exc)) from exc
return downloads.settings_view()
@router.post("/pause-all")
def pause_all_downloads():
return downloads.pause_all()
@router.post("/resume-all")
def resume_all_downloads():
return downloads.resume_all()
@router.post("/{drama_id}/pause")
def pause_download(drama_id: str):
return downloads.pause(drama_id)
@router.post("/{drama_id}/resume")
def resume_download(drama_id: str):
return downloads.resume(drama_id)
@router.post("/{drama_id}/cancel")
def cancel_download(drama_id: str):
return downloads.cancel(drama_id)
@router.post("/{drama_id}/open")
def open_download(drama_id: str):
def action():
downloads.open_folder(drama_id)
return downloads.snapshot()
return call(action)
@router.delete("/{drama_id}")
def remove_download(drama_id: str):
return downloads.remove(drama_id)
+33
View File
@@ -0,0 +1,33 @@
import asyncio
from fastapi import APIRouter, WebSocket, WebSocketDisconnect
from server.services.downloads import downloads
from server.services.live import live
from server.services.sessions import sessions
router = APIRouter()
def _payload() -> dict:
profile = sessions.snapshot()
return {
"downloads": downloads.snapshot(),
"profile": profile["profile"] if profile else None,
}
@router.websocket("/api/live")
async def live_socket(websocket: WebSocket):
await websocket.accept()
seen = -1
try:
while True:
revision = await asyncio.to_thread(live.wait, seen, 20)
if revision != seen:
seen = revision
await websocket.send_json(_payload())
else:
await websocket.send_json({"type": "ping"})
except WebSocketDisconnect:
return
+89
View File
@@ -0,0 +1,89 @@
from urllib.parse import urlsplit
import httpx
from fastapi import APIRouter, HTTPException
from fastapi.responses import Response, StreamingResponse
from pydantic import BaseModel
from server.schemas import DeviceBody, device_payload
from server.services.device import DeviceProfile
from server.services.http import http
from server.services import playlist as playlists
from server.services.sessions import sessions
router = APIRouter(prefix="/api/play")
class PlayBody(BaseModel):
token: str = ""
device: DeviceBody | None = None
vid: int
def _device(token: str, raw: dict | None) -> DeviceProfile | None:
if not token:
return None
if raw:
return sessions.bind(token, DeviceProfile.merge(raw))
return sessions.device_for(token)
def _public_url(url: str) -> str:
parts = urlsplit(url)
if parts.scheme not in ("http", "https") or not parts.hostname:
raise HTTPException(status_code=400, detail="播放地址无效")
return url
@router.post("/episodes")
def list_episodes(body: PlayBody):
device = _device(body.token, device_payload(body.device))
try:
rows = playlists.load_episodes(body.token, device, body.vid)
except RuntimeError as exc:
raise HTTPException(status_code=400, detail=str(exc)) from exc
return {
"episodes": [
{
"cid": ep["cid"],
"idx": ep["idx"],
"title": ep["title"],
"playable": bool(ep.get("src")),
}
for ep in rows
]
}
@router.get("/segment")
def play_segment(url: str):
target = _public_url(url)
try:
upstream = http.open_stream(target)
except (httpx.HTTPError, TimeoutError, OSError, RuntimeError) as exc:
raise HTTPException(status_code=502, detail="分片读取失败") from exc
media = upstream.headers.get("content-type", "video/mp2t").split(";", 1)[0].strip() or "video/mp2t"
def chunks():
try:
for chunk in upstream.iter_bytes(64 * 1024):
if chunk:
yield chunk
finally:
upstream.close()
return StreamingResponse(chunks(), media_type=media)
@router.get("/{vid}/{cid}/index.m3u8")
def episode_playlist(vid: int, cid: int):
token = sessions.active_token()
device = sessions.device_for(token) if token else None
src = playlists.episode_src(token, device, vid, cid)
if not src:
raise HTTPException(status_code=404, detail="这一集还没有可播放的地址")
try:
text = playlists.playback_playlist(src)
except RuntimeError as exc:
raise HTTPException(status_code=502, detail=str(exc)) from exc
return Response(text, media_type="application/vnd.apple.mpegurl", headers={"Cache-Control": "no-store"})
+668
View File
@@ -0,0 +1,668 @@
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()
+75
View File
@@ -0,0 +1,75 @@
import gzip
import shutil
import subprocess
from pathlib import Path
import httpx
FFMPEG_DIR = Path(__file__).resolve().parents[2] / "tools" / "ffmpeg"
FFMPEG_EXE = FFMPEG_DIR / "ffmpeg.exe"
FFMPEG_URL = "https://github.com/eugeneware/ffmpeg-static/releases/download/b6.1.1/ffmpeg-win32-x64.gz"
def ffmpeg_path() -> Path | None:
if FFMPEG_EXE.is_file():
return FFMPEG_EXE
found = shutil.which("ffmpeg")
return Path(found) if found else None
def ensure_ffmpeg() -> Path:
current = ffmpeg_path()
if current is not None:
return current
FFMPEG_DIR.mkdir(parents=True, exist_ok=True)
archive = FFMPEG_DIR / "ffmpeg-win32-x64.gz"
try:
with httpx.Client(follow_redirects=True, timeout=120.0, trust_env=False) as client:
response = client.get(FFMPEG_URL)
response.raise_for_status()
archive.write_bytes(response.content)
import gzip
FFMPEG_EXE.write_bytes(gzip.decompress(archive.read_bytes()))
finally:
archive.unlink(missing_ok=True)
if not FFMPEG_EXE.is_file():
raise RuntimeError("没有找到 ffmpeg")
return FFMPEG_EXE
def convert_to_mp4(episode_dir: Path, mp4_path: Path) -> None:
exe = ensure_ffmpeg()
mp4_path.parent.mkdir(parents=True, exist_ok=True)
temporary = mp4_path.with_name(f"{mp4_path.stem}.part.mp4")
completed = subprocess.run(
[
str(exe),
"-y",
"-allowed_extensions",
"ALL",
"-i",
"video.m3u8",
"-c",
"copy",
"-movflags",
"+faststart",
"-f",
"mp4",
str(temporary),
],
cwd=episode_dir,
capture_output=True,
check=False,
)
if completed.returncode != 0 or not temporary.is_file() or temporary.stat().st_size <= 0:
temporary.unlink(missing_ok=True)
detail = (completed.stderr or b"").decode("utf-8", "replace")[-400:]
raise RuntimeError("转换成 MP4 失败" + (f":{detail}" if detail else ""))
temporary.replace(mp4_path)
def unused_zip_helper(path: Path) -> None:
if path.suffix == ".zip":
with zipfile.ZipFile(path) as archive:
archive.extractall(path.parent)
+58
View File
@@ -89,6 +89,64 @@ class HttpClient:
raise RuntimeError("不是图片")
return response.content, media
def download_file(self, url: str, dest, on_chunk=None, stop=None) -> int:
from pathlib import Path
target = Path(dest)
target.parent.mkdir(parents=True, exist_ok=True)
temporary = target.with_name(target.name + ".part")
last_error = None
for _attempt in range(3):
if stop and stop():
raise RuntimeError("已取消")
try:
with self._client_locked().stream("GET", url, headers={"User-Agent": "RaptDrama"}) as response:
response.raise_for_status()
size = 0
with temporary.open("wb") as handle:
for chunk in response.iter_bytes(64 * 1024):
if stop and stop():
raise RuntimeError("已取消")
if not chunk:
continue
handle.write(chunk)
size += len(chunk)
if on_chunk:
on_chunk(len(chunk))
temporary.replace(target)
return size
except RuntimeError:
temporary.unlink(missing_ok=True)
raise
except (httpx.HTTPError, TimeoutError, OSError) as exc:
last_error = exc
temporary.unlink(missing_ok=True)
time.sleep(0.4)
raise RuntimeError(f"下载失败: {last_error}")
def read_text(self, url: str) -> str:
last_error = None
for _attempt in range(3):
try:
response = self._client_locked().get(url, headers={"User-Agent": "RaptDrama"})
response.raise_for_status()
return response.text
except (httpx.HTTPError, TimeoutError, OSError) as exc:
last_error = exc
time.sleep(0.4)
raise RuntimeError(f"下载失败: {last_error}")
def open_stream(self, url: str) -> httpx.Response:
client = self._client_locked()
request = client.build_request("GET", url, headers={"User-Agent": "RaptDrama"})
response = client.send(request, stream=True)
try:
response.raise_for_status()
except Exception:
response.close()
raise
return response
def _close_locked(self) -> None:
if self._client is not None:
self._client.close()
+21
View File
@@ -0,0 +1,21 @@
import threading
class LiveHub:
def __init__(self) -> None:
self._condition = threading.Condition()
self.revision = 0
def touch(self) -> None:
with self._condition:
self.revision += 1
self._condition.notify_all()
def wait(self, seen: int, timeout: float = 20) -> int:
with self._condition:
if self.revision == seen:
self._condition.wait(timeout)
return self.revision
live = LiveHub()
+218
View File
@@ -0,0 +1,218 @@
import json
import re
from pathlib import Path
from urllib.parse import quote, urljoin
from server.services.device import DeviceProfile
from server.services.http import http
from server.services.rapt import API_BASE
CACHE_DIR = Path(__file__).resolve().parents[2] / "cache" / "playlists"
def safe_name(text: str, fallback: str) -> str:
cleaned = re.sub(r'[<>:"/\\|?*\x00-\x1f]', " ", text or "")
cleaned = re.sub(r"\s+", " ", cleaned).strip().rstrip(". ")
return (cleaned[:80] or fallback)
def _items(data):
if isinstance(data, list):
return data, len(data)
items = (data or {}).get("list") or []
return items, (data or {}).get("count")
def _normalize(item: dict) -> dict | None:
cid = item.get("cid") or item.get("id")
if cid is None:
return None
idx = item.get("idx") or item.get("number") or 0
title = str(item.get("msg") or item.get("title") or f"EP.{idx}").strip()
return {
"cid": int(cid),
"idx": int(idx) if str(idx).isdigit() else 0,
"title": title,
"src": str(item.get("src") or "").strip(),
}
def _request_playlist(token: str, device: DeviceProfile, vid: int) -> list[dict]:
try:
data = http.request_data("POST", API_BASE + "/app/video/playlist", token, device, {
"vid": vid,
"page": 1,
"before": 1,
})
items, total = _items(data)
episodes = [row for row in (_normalize(item) for item in items if isinstance(item, dict)) if row]
if episodes and (total is None or len(episodes) >= int(total)):
return episodes
except RuntimeError:
episodes = []
collected = []
total = None
for page in range(1, 41):
data = http.request_data("POST", API_BASE + "/app/video/playlistv2", token, device, {
"vid": vid,
"page": page,
"before": 1 if page == 1 else 0,
})
items, page_total = _items(data)
if page_total is not None:
total = page_total
if not items:
break
for item in items:
if isinstance(item, dict):
row = _normalize(item)
if row:
collected.append(row)
if total is not None and len(collected) >= int(total):
break
return collected or episodes
def _cache_path(vid: int) -> Path:
return CACHE_DIR / f"{vid}.json"
def _read_cache(vid: int) -> list[dict]:
path = _cache_path(vid)
if not path.is_file():
return []
try:
data = json.loads(path.read_text(encoding="utf-8"))
except (OSError, json.JSONDecodeError):
return []
rows = data.get("episodes") if isinstance(data, dict) else None
return [row for row in rows if isinstance(row, dict)] if isinstance(rows, list) else []
def _write_cache(vid: int, episodes: list[dict]) -> None:
CACHE_DIR.mkdir(parents=True, exist_ok=True)
payload = {
"episodes": [
{"cid": ep["cid"], "idx": ep["idx"], "title": ep["title"], "src": ep.get("src") or ""}
for ep in episodes
]
}
path = _cache_path(vid)
temporary = path.with_suffix(".json.tmp")
temporary.write_text(json.dumps(payload, ensure_ascii=False), encoding="utf-8")
temporary.replace(path)
def _filled(episodes: list[dict]) -> int:
return sum(1 for ep in episodes if ep.get("src"))
def merge_episodes(cached: list[dict], fresh: list[dict]) -> tuple[list[dict], bool]:
by_cid: dict[int, dict] = {}
for source in (cached, fresh):
for ep in source:
cid = int(ep["cid"])
current = by_cid.get(cid)
if current is None:
by_cid[cid] = {
"cid": cid,
"idx": int(ep.get("idx") or 0),
"title": ep.get("title") or f"EP.{ep.get('idx') or cid}",
"src": ep.get("src") or "",
}
continue
if ep.get("title"):
current["title"] = ep["title"]
if ep.get("idx"):
current["idx"] = int(ep["idx"])
if ep.get("src"):
current["src"] = ep["src"]
merged = sorted(by_cid.values(), key=lambda ep: (ep["idx"], ep["cid"]))
return merged, _filled(fresh) >= _filled(cached) and _filled(fresh) > 0
def load_episodes(token: str, device: DeviceProfile | None, vid: int) -> list[dict]:
cached = _read_cache(vid)
fresh = []
if token and device is not None:
try:
fresh = _request_playlist(token, device, vid)
except RuntimeError:
fresh = []
merged, save = merge_episodes(cached, fresh)
if save:
_write_cache(vid, merged)
elif not cached and merged:
_write_cache(vid, merged)
return merged
def expand_variant(master_url: str) -> tuple[str, str, list[str]]:
master_text = http.read_text(master_url)
variant_rel = next(
(line.strip() for line in master_text.splitlines() if line.strip() and not line.startswith("#")),
"",
)
if not variant_rel:
raise RuntimeError("播放列表没有视频地址")
variant_url = urljoin(master_url, variant_rel)
media_text = http.read_text(variant_url)
segments = [
urljoin(variant_url, line.strip())
for line in media_text.splitlines()
if line.strip() and not line.startswith("#")
]
if not segments:
raise RuntimeError("播放列表没有分片")
return master_text, media_text, segments
def episode_src(token: str, device: DeviceProfile | None, vid: int, cid: int) -> str:
for ep in load_episodes(token, device, vid):
if int(ep.get("cid") or 0) == cid and ep.get("src"):
return str(ep["src"])
return ""
def playback_playlist(master_url: str) -> str:
master_text = http.read_text(master_url)
variant_rel = next(
(line.strip() for line in master_text.splitlines() if line.strip() and not line.startswith("#")),
"",
)
if not variant_rel:
raise RuntimeError("播放列表没有视频地址")
if _looks_like_segment(variant_rel):
return _rewrite_playlist(master_text, master_url)
variant_url = urljoin(master_url, variant_rel)
return _rewrite_playlist(http.read_text(variant_url), variant_url)
def _looks_like_segment(url: str) -> bool:
path = url.split("?", 1)[0].lower()
return path.endswith((".ts", ".m4s", ".m4v", ".mp4", ".aac", ".vtt"))
def _proxy_url(url: str) -> str:
return "/api/play/segment?url=" + quote(url, safe="")
def _rewrite_playlist(text: str, base: str) -> str:
lines = []
for line in text.splitlines():
stripped = line.strip()
if not stripped:
lines.append(line)
continue
if stripped.startswith("#"):
lines.append(_rewrite_tag(line, base))
continue
lines.append(_proxy_url(urljoin(base, stripped)))
return "\n".join(lines) + "\n"
def _rewrite_tag(line: str, base: str) -> str:
def replace(match: re.Match) -> str:
return 'URI="' + _proxy_url(urljoin(base, match.group(1))) + '"'
return re.sub(r'URI="([^"]*)"', replace, line)
+4
View File
@@ -1,6 +1,7 @@
import threading
from server.services.device import DeviceProfile
from server.services.live import live
class SessionStore:
@@ -26,6 +27,7 @@ class SessionStore:
self._active = token
self._profile = dict(profile)
self._revision += 1
live.touch()
def set_profile(self, token: str, profile: dict) -> bool:
with self._lock:
@@ -33,6 +35,7 @@ class SessionStore:
return False
self._profile = dict(profile)
self._revision += 1
live.touch()
return True
def active_token(self) -> str:
@@ -51,6 +54,7 @@ class SessionStore:
if self._active == token:
self._active = ""
self._profile = None
live.touch()
sessions = SessionStore()