This commit is contained in:
root
2026-09-26 13:55:43 +08:00
commit 9dfc684869
40 changed files with 5632 additions and 0 deletions
+3
View File
@@ -0,0 +1,3 @@
__pycache__/
.pytest_cache/
*.py[cod]
+35
View File
@@ -0,0 +1,35 @@
# m3u8downloaderd
`m3u8downloaderd` is an in-memory local-network web interface for the adjacent
`N_m3u8DL_py` package. It requires Python 3.10+ and a system `ffmpeg`.
Install both local packages into the same environment:
```bash
python -m pip install -e ./N_m3u8DL_py -e ./m3u8downloaderd
```
Start the service:
```bash
python -m m3u8downloaderd \
--host 0.0.0.0 \
--port 8000 \
--download-root /tmp \
--max-active-downloads 2
```
The service exposes no authentication. Use it only on a trusted network. The
folder chooser is restricted to `--download-root` and its descendants.
Tasks, progress, history, and settings exist only while the process is running;
restarting the service clears them. Downloaded media files remain on disk.
`--max-active-downloads` sets the initial video concurrency. It can be changed
while the service is running from the Settings tab, with a supported range of
1 to 32 videos. Runtime logs are available in the Log tab and remain in memory
only.
Run tests from the repository root:
```bash
PYTHONPATH=N_m3u8DL_py/src:m3u8downloaderd/src python -m pytest m3u8downloaderd/tests
```
+27
View File
@@ -0,0 +1,27 @@
[build-system]
requires = ["setuptools>=68"]
build-backend = "setuptools.build_meta"
[project]
name = "m3u8downloaderd"
version = "0.1.0"
description = "Local-network web service for N_m3u8DL-PY downloads"
requires-python = ">=3.10"
dependencies = ["n-m3u8dl-py>=0.2.0"]
[project.scripts]
m3u8downloaderd = "m3u8downloaderd.__main__:main"
[tool.setuptools]
package-dir = {"" = "src"}
include-package-data = true
[tool.setuptools.packages.find]
where = ["src"]
[tool.setuptools.package-data]
m3u8downloaderd = ["static/*"]
[tool.pytest.ini_options]
testpaths = ["tests"]
addopts = "-q"
@@ -0,0 +1,3 @@
"""Durable web front end for N_m3u8DL-PY."""
__version__ = "0.1.0"
@@ -0,0 +1,16 @@
from __future__ import annotations
from .web import build_argument_parser, config_from_args, serve
def main() -> None:
parser = build_argument_parser()
args = parser.parse_args()
try:
serve(config_from_args(args))
except (RuntimeError, ValueError) as error:
parser.error(str(error))
if __name__ == "__main__":
main()
@@ -0,0 +1,156 @@
"""Input validation and path rules shared by the service layers."""
from __future__ import annotations
import base64
import binascii
import json
import re
from dataclasses import dataclass
from datetime import date, datetime
from pathlib import Path
from urllib.parse import urlsplit
REQUIRED_FIELDS = ("code", "title", "href", "image_src", "m3u8_url", "m3u8_referer")
DATE_TEMPLATE = "YYYY-MM-DD_HH-mm-ss"
LEGACY_DATE_TEMPLATE = "YYYY-MM-DD"
DEFAULT_PROXY = "socks5://192.168.4.100:10808"
MAX_ACTIVE_DOWNLOADS = 32
_SEPARATORS = re.compile(r"[\\/\x00]")
class ValidationError(ValueError):
"""Raised when a request cannot become a runnable task."""
@dataclass(frozen=True)
class DownloadItemInput:
code: object
title: str
href: object
image_src: str
m3u8_url: str
m3u8_referer: str
def resolve_default_title(template: str, now: datetime | date | None = None) -> str:
"""Expand date templates; all other values are literal."""
if template == DATE_TEMPLATE:
current = now or datetime.now()
if isinstance(current, datetime):
return current.strftime("%Y-%m-%d_%H-%M-%S")
return f"{current.isoformat()}_00-00-00"
if template == LEGACY_DATE_TEMPLATE:
return (now or date.today()).isoformat()
return template
def decode_download_items(value: str) -> list[DownloadItemInput]:
"""Accept JSON or standard/URL-safe base64 that decodes to UTF-8 JSON."""
if not isinstance(value, str) or not value.strip():
raise ValidationError("Download information is required")
payload = value.strip()
try:
decoded = json.loads(payload)
except json.JSONDecodeError:
decoded = _decode_base64_json(payload)
if not isinstance(decoded, list) or not decoded:
raise ValidationError("Download information must be a non-empty JSON list")
items: list[DownloadItemInput] = []
for index, item in enumerate(decoded, 1):
if not isinstance(item, dict):
raise ValidationError(f"Item {index} must be a JSON object")
missing = [field for field in REQUIRED_FIELDS if field not in item]
if missing:
raise ValidationError(f"Item {index} is missing: {', '.join(missing)}")
title = _nonempty_string(item["title"], index, "title")
image_src = _http_url(item["image_src"], index, "image_src")
m3u8_url = _http_url(item["m3u8_url"], index, "m3u8_url")
referer = _nonempty_string(item["m3u8_referer"], index, "m3u8_referer")
items.append(
DownloadItemInput(
code=item["code"],
title=title,
href=item["href"],
image_src=image_src,
m3u8_url=m3u8_url,
m3u8_referer=referer,
)
)
return items
def validate_proxy(value: object) -> str:
if value in (None, ""):
return ""
if not isinstance(value, str):
raise ValidationError("Proxy must be a string")
parsed = urlsplit(value.strip())
if parsed.scheme not in {"http", "https", "socks5", "socks5h"} or not parsed.netloc:
raise ValidationError("Proxy must be an HTTP(S) or SOCKS5 URL")
if parsed.username is not None or parsed.password is not None:
raise ValidationError("Authenticated proxies are not supported")
return value.strip()
def validate_max_active_downloads(value: object) -> int:
if isinstance(value, bool):
raise ValidationError("Concurrent downloads must be an integer")
if isinstance(value, str):
value = value.strip()
if not value.isdecimal():
raise ValidationError("Concurrent downloads must be an integer")
value = int(value)
if not isinstance(value, int):
raise ValidationError("Concurrent downloads must be an integer")
if not 1 <= value <= MAX_ACTIVE_DOWNLOADS:
raise ValidationError(f"Concurrent downloads must be between 1 and {MAX_ACTIVE_DOWNLOADS}")
return value
def validate_folder_title(value: object) -> str:
if not isinstance(value, str):
raise ValidationError("Task title must be a string")
value = value.strip()
if not value:
return ""
if value in {".", ".."} or _SEPARATORS.search(value):
raise ValidationError("Task title must be one directory name")
return value
def ensure_within_root(value: str | Path, root: Path, *, must_exist: bool = True) -> Path:
root = root.resolve()
candidate = Path(value).expanduser().resolve(strict=False)
try:
candidate.relative_to(root)
except ValueError as error:
raise ValidationError("Directory must be inside the configured download root") from error
if must_exist and (not candidate.exists() or not candidate.is_dir()):
raise ValidationError("Selected directory does not exist")
return candidate
def _decode_base64_json(value: str) -> object:
compact = "".join(value.split())
padded = compact + "=" * (-len(compact) % 4)
try:
raw = base64.b64decode(padded, altchars=b"-_", validate=True)
return json.loads(raw.decode("utf-8"))
except (binascii.Error, UnicodeDecodeError, json.JSONDecodeError) as error:
raise ValidationError("Download information must be JSON or base64-encoded JSON") from error
def _nonempty_string(value: object, index: int, field: str) -> str:
if not isinstance(value, str) or not value.strip():
raise ValidationError(f"Item {index} field {field} must be a non-empty string")
return value.strip()
def _http_url(value: object, index: int, field: str) -> str:
result = _nonempty_string(value, index, field)
parsed = urlsplit(result)
if parsed.scheme not in {"http", "https"} or not parsed.netloc:
raise ValidationError(f"Item {index} field {field} must be an HTTP(S) URL")
return result
@@ -0,0 +1,340 @@
const state = {
bootstrap: null,
tasks: [],
directoryTarget: null,
currentDirectory: null,
validationTimer: null,
openActivities: new Set(),
activeTab: "downloads",
logs: [],
logLevel: "",
};
const $ = (selector) => document.querySelector(selector);
async function api(path, options = {}) {
const response = await fetch(path, {
headers: { "Content-Type": "application/json", ...(options.headers || {}) },
...options,
});
const body = await response.json().catch(() => ({}));
if (!response.ok) throw new Error(body.error || `Request failed (${response.status})`);
return body;
}
function text(value) {
return value == null || value === "" ? "-" : String(value);
}
function label(value) {
return String(value || "unknown").replaceAll("_", " ");
}
function make(tag, className, content) {
const node = document.createElement(tag);
if (className) node.className = className;
if (content != null) node.textContent = content;
return node;
}
async function refreshBootstrap() {
state.bootstrap = await api("/api/bootstrap");
state.tasks = state.bootstrap.tasks;
renderTasks();
if (state.activeTab === "logs") await refreshLogs();
}
async function refreshLogs() {
const query = state.logLevel ? `?level=${encodeURIComponent(state.logLevel)}` : "";
const result = await api(`/api/logs${query}`);
state.logs = result.logs;
renderLogs();
}
function renderLogs() {
const list = $("#log-list");
const empty = $("#empty-logs");
list.replaceChildren();
empty.hidden = state.logs.length !== 0;
for (const entry of state.logs) {
const row = make("li", `log-row log-${entry.level}`);
const metadata = [entry.created_at.replace("T", " ").replace("+00:00", " UTC"), entry.level.toUpperCase()];
if (entry.task_id != null) metadata.push(`Task ${entry.task_id}`);
if (entry.item_id != null) metadata.push(`Video ${entry.item_id}`);
row.append(make("span", "log-meta", metadata.join(" · ")));
row.append(make("span", "log-message", entry.message));
list.append(row);
}
}
function renderTasks() {
const list = $("#task-list");
const empty = $("#empty-state");
list.querySelectorAll(".task-row[data-task-id] .event-log[open]").forEach((entry) => {
state.openActivities.add(entry.closest(".task-row").dataset.taskId);
});
list.replaceChildren();
$("#queue-count").textContent = `${state.tasks.length} ${state.tasks.length === 1 ? "task" : "tasks"}`;
empty.hidden = state.tasks.length !== 0;
for (const task of state.tasks) list.append(renderTask(task));
}
function renderTask(task) {
const fragment = $("#task-template").content.cloneNode(true);
const row = fragment.querySelector(".task-row");
row.dataset.status = task.status;
row.dataset.taskId = String(task.id);
const main = fragment.querySelector(".task-main");
main.append(make("div", "task-title", task.title || task.base_dir));
main.append(make("div", "task-meta", task.output_dir));
main.append(make("span", "task-status", label(task.status)));
const actions = fragment.querySelector(".task-actions");
if (["queued", "running"].includes(task.status)) {
const cancel = make("button", "secondary-action", "Cancel");
cancel.type = "button";
cancel.addEventListener("click", () => mutate(`/api/tasks/${task.id}/cancel`, "Cancel this task?"));
actions.append(cancel);
}
if (["partial", "failed", "cancelled"].includes(task.status)) {
const retry = make("button", "secondary-action", "Retry task");
retry.type = "button";
retry.addEventListener("click", () => mutate(`/api/tasks/${task.id}/retry`));
actions.append(retry);
}
const progress = taskProgress(task.items);
const progressRoot = fragment.querySelector(".task-progress");
const caption = make("div", "progress-caption");
caption.append(make("span", "", progress.caption));
caption.append(make("span", "", `${progress.percent}%`));
const meter = document.createElement("progress");
meter.max = 100;
meter.value = progress.percent;
progressRoot.append(caption, meter);
const items = fragment.querySelector(".item-list");
for (const item of task.items) items.append(renderItem(item));
const eventLog = fragment.querySelector(".event-log");
const eventList = eventLog.querySelector("ol");
if (!task.events.length) {
state.openActivities.delete(String(task.id));
eventLog.remove();
}
else {
eventLog.open = state.openActivities.has(String(task.id));
eventLog.addEventListener("toggle", () => {
if (eventLog.open) state.openActivities.add(String(task.id));
else state.openActivities.delete(String(task.id));
});
for (const event of task.events) {
const entry = make("li", "", `${event.created_at.replace("T", " ").replace("+00:00", " UTC")} · ${event.message}`);
eventList.append(entry);
}
}
return fragment;
}
function renderItem(item) {
const row = make("div", "item-row");
const details = make("div", "");
details.append(make("div", "item-name", item.title));
details.append(make("div", "item-detail", item.stage));
const progress = itemProgress(item);
const progressRoot = make("div", "item-progress");
const caption = make("div", "progress-caption");
caption.append(make("span", "", progress.caption));
caption.append(make("span", "", progress.total ? `${progress.percent}%` : "Waiting"));
const meter = document.createElement("progress");
meter.max = 100;
if (progress.total) meter.value = progress.percent;
progressRoot.append(caption, meter);
details.append(progressRoot);
if (item.warning) details.append(make("div", "item-detail item-warning", item.warning));
if (item.error) details.append(make("div", "item-detail item-warning", item.error));
row.append(details);
const actions = make("div", "item-actions");
if (["failed", "cancelled", "completed_warning"].includes(item.status)) {
const retry = make("button", "secondary-action", item.status === "completed_warning" ? "Retry cover" : "Retry video");
retry.type = "button";
retry.addEventListener("click", () => mutate(`/api/items/${item.id}/retry`));
actions.append(retry);
}
if (item.output_path) actions.append(make("span", "item-detail", "MP4 ready"));
row.append(actions);
return row;
}
function taskProgress(items) {
const completed = items.filter((item) => ["completed", "completed_warning"].includes(item.status)).length;
const total = items.length;
return {
percent: Math.round((completed / Math.max(total, 1)) * 100),
caption: `${completed}/${total} videos completed`,
};
}
function itemProgress(item) {
const total = item.total_segments || 0;
const completed = item.completed_segments || 0;
if (total) {
return {
total,
percent: Math.min(100, Math.round((completed / total) * 100)),
caption: `${completed}/${total} segments`,
};
}
if (["completed", "completed_warning"].includes(item.status)) {
return { total: 1, percent: 100, caption: "All segments downloaded" };
}
return { total: 0, percent: 0, caption: "Waiting for playlist" };
}
async function mutate(path, confirmation) {
if (confirmation && !window.confirm(confirmation)) return;
try {
await api(path, { method: "POST", body: "{}" });
await refreshBootstrap();
} catch (error) {
window.alert(error.message);
}
}
function showTaskDialog() {
const modal = $("#task-dialog");
const settings = state.bootstrap.settings;
$("#task-title").value = state.bootstrap.default_title;
$("#task-directory").value = settings.default_directory;
$("#task-payload").value = "";
setValidation("", false);
modal.showModal();
}
function setValidation(message, valid) {
const status = $("#payload-status");
status.textContent = message;
status.classList.toggle("is-valid", valid);
status.classList.toggle("is-invalid", Boolean(message) && !valid);
$("#start-task").disabled = !valid;
}
async function validateTaskPayload() {
const payload = $("#task-payload").value;
if (!payload.trim()) return setValidation("", false);
try {
const result = await api("/api/validate", { method: "POST", body: JSON.stringify({ payload }) });
setValidation(`${result.count} ${result.count === 1 ? "video" : "videos"} ready`, true);
} catch (error) {
setValidation(error.message, false);
}
}
function openDirectoryPicker(target) {
state.directoryTarget = target;
const input = target === "task" ? $("#task-directory") : $("#setting-directory");
loadDirectory(input.value).then(() => $("#directory-dialog").showModal()).catch((error) => window.alert(error.message));
}
async function loadDirectory(path) {
const result = await api(`/api/directories?path=${encodeURIComponent(path)}`);
state.currentDirectory = result;
$("#directory-path").textContent = result.path;
$("#directory-up").disabled = !result.parent;
const list = $("#directory-list");
list.replaceChildren();
for (const directory of result.directories) {
const button = make("button", "directory-entry", directory.name);
button.type = "button";
button.addEventListener("click", () => loadDirectory(directory.path));
list.append(button);
}
if (!result.directories.length) list.append(make("p", "item-detail", "No subfolders."));
}
async function submitTask(event) {
event.preventDefault();
const start = $("#start-task");
start.disabled = true;
try {
await api("/api/tasks", {
method: "POST",
body: JSON.stringify({
title: $("#task-title").value,
payload: $("#task-payload").value,
directory: $("#task-directory").value,
}),
});
$("#task-dialog").close();
await refreshBootstrap();
} catch (error) {
setValidation(error.message, false);
}
}
async function submitSettings(event) {
event.preventDefault();
const status = $("#settings-status");
try {
const settings = await api("/api/settings", {
method: "PUT",
body: JSON.stringify({
default_title_template: $("#setting-title").value,
default_directory: $("#setting-directory").value,
proxy: $("#setting-proxy").value,
max_active_downloads: Number($("#setting-max-downloads").value),
}),
});
state.bootstrap.settings = settings;
await refreshBootstrap();
status.textContent = "Saved";
} catch (error) {
status.textContent = error.message;
}
}
function switchTab(tab) {
state.activeTab = tab;
document.querySelectorAll(".tab").forEach((button) => button.classList.toggle("is-active", button.dataset.tab === tab));
document.querySelectorAll(".view").forEach((view) => view.classList.toggle("is-active", view.id === `${tab}-view`));
if (tab === "logs") refreshLogs().catch((error) => window.alert(error.message));
}
function bind() {
$("#add-task").addEventListener("click", showTaskDialog);
document.querySelectorAll(".tab").forEach((button) => button.addEventListener("click", () => switchTab(button.dataset.tab)));
$("#close-task").addEventListener("click", () => $("#task-dialog").close());
$("#cancel-task").addEventListener("click", () => $("#task-dialog").close());
$("#task-form").addEventListener("submit", submitTask);
$("#settings-form").addEventListener("submit", submitSettings);
$("#task-payload").addEventListener("input", () => {
window.clearTimeout(state.validationTimer);
state.validationTimer = window.setTimeout(validateTaskPayload, 250);
});
$("#browse-task").addEventListener("click", () => openDirectoryPicker("task"));
$("#browse-settings").addEventListener("click", () => openDirectoryPicker("settings"));
$("#log-level").addEventListener("change", () => {
state.logLevel = $("#log-level").value;
refreshLogs().catch((error) => window.alert(error.message));
});
$("#close-directory").addEventListener("click", () => $("#directory-dialog").close());
$("#directory-up").addEventListener("click", () => loadDirectory(state.currentDirectory.parent));
$("#select-directory").addEventListener("click", () => {
const selector = state.directoryTarget === "task" ? "#task-directory" : "#setting-directory";
$(selector).value = state.currentDirectory.path;
$("#directory-dialog").close();
});
}
async function initialize() {
bind();
await refreshBootstrap();
const settings = state.bootstrap.settings;
$("#setting-title").value = settings.default_title_template;
$("#setting-directory").value = settings.default_directory;
$("#setting-proxy").value = settings.proxy;
$("#setting-max-downloads").value = settings.max_active_downloads;
window.setInterval(() => refreshBootstrap().catch(() => {}), 1000);
}
initialize().catch((error) => window.alert(error.message));
@@ -0,0 +1,154 @@
<!doctype html>
<html lang="en">
<head>
<meta charset="utf-8">
<meta name="viewport" content="width=device-width, initial-scale=1">
<title>m3u8downloaderd</title>
<link rel="stylesheet" href="/static/styles.css">
</head>
<body>
<header class="topbar">
<div class="brand" aria-label="m3u8downloaderd">
<span class="brand-mark" aria-hidden="true"><i></i><i></i><i></i></span>
<span>m3u8downloaderd</span>
</div>
<nav class="tabs" aria-label="Application">
<button class="tab is-active" type="button" data-tab="downloads">Downloads</button>
<button class="tab" type="button" data-tab="settings">Settings</button>
<button class="tab" type="button" data-tab="logs">Log</button>
</nav>
<button class="primary-action" id="add-task" type="button">Add Download Task</button>
</header>
<main>
<section id="downloads-view" class="view is-active" aria-labelledby="downloads-heading">
<div class="view-heading">
<div>
<p class="eyebrow">Queue</p>
<h1 id="downloads-heading">Downloads</h1>
</div>
<span class="queue-count" id="queue-count">0 tasks</span>
</div>
<div id="task-list" class="task-list" aria-live="polite"></div>
<p id="empty-state" class="empty-state">No download tasks.</p>
</section>
<section id="settings-view" class="view" aria-labelledby="settings-heading">
<div class="view-heading">
<div>
<p class="eyebrow">Service</p>
<h1 id="settings-heading">Settings</h1>
</div>
</div>
<form id="settings-form" class="settings-form">
<label>
<span>Default title</span>
<input id="setting-title" name="default_title_template" autocomplete="off">
</label>
<label>
<span>Default folder</span>
<span class="path-control">
<input id="setting-directory" name="default_directory" readonly>
<button class="icon-button" id="browse-settings" type="button" aria-label="Choose default folder" title="Choose default folder">Folder</button>
</span>
</label>
<label>
<span>Proxy</span>
<input id="setting-proxy" name="proxy" inputmode="url" placeholder="socks5://192.168.4.100:10808" autocomplete="off">
</label>
<label>
<span>Concurrent video downloads</span>
<input id="setting-max-downloads" name="max_active_downloads" type="number" min="1" max="32" step="1" required>
</label>
<div class="form-actions">
<button class="primary-action" type="submit">Save settings</button>
<span id="settings-status" class="form-status" role="status"></span>
</div>
</form>
</section>
<section id="logs-view" class="view" aria-labelledby="logs-heading">
<div class="view-heading">
<div>
<p class="eyebrow">Runtime</p>
<h1 id="logs-heading">Log</h1>
</div>
<label class="log-filter">
<span>Level</span>
<select id="log-level">
<option value="">All</option>
<option value="debug">Debug</option>
<option value="info">Info</option>
<option value="warning">Warning</option>
<option value="error">Error</option>
</select>
</label>
</div>
<ol id="log-list" class="log-list" aria-live="polite"></ol>
<p id="empty-logs" class="empty-state" hidden>No log entries.</p>
</section>
</main>
<dialog id="task-dialog" class="modal task-modal">
<form id="task-form" method="dialog">
<header class="modal-header">
<h2>Add Download Task</h2>
<button class="close-button" id="close-task" type="button" aria-label="Close" title="Close">Close</button>
</header>
<label>
<span>Title</span>
<input id="task-title" name="title" autocomplete="off">
</label>
<label class="payload-label">
<span>Download information</span>
<textarea id="task-payload" name="payload" spellcheck="false"></textarea>
</label>
<label>
<span>Folder</span>
<span class="path-control">
<input id="task-directory" name="directory" readonly>
<button class="icon-button" id="browse-task" type="button" aria-label="Choose folder" title="Choose folder">Folder</button>
</span>
</label>
<p id="payload-status" class="validation-status" role="status"></p>
<footer class="modal-footer">
<button class="secondary-action" id="cancel-task" type="button">Cancel</button>
<button class="primary-action" id="start-task" type="submit" disabled>Start</button>
</footer>
</form>
</dialog>
<dialog id="directory-dialog" class="modal directory-modal">
<header class="modal-header">
<div>
<p class="eyebrow">Server folder</p>
<h2>Choose folder</h2>
</div>
<button class="close-button" id="close-directory" type="button" aria-label="Close" title="Close">Close</button>
</header>
<p id="directory-path" class="directory-path"></p>
<div class="directory-actions">
<button class="secondary-action" id="directory-up" type="button">Up</button>
<button class="primary-action" id="select-directory" type="button">Select folder</button>
</div>
<div id="directory-list" class="directory-list"></div>
</dialog>
<template id="task-template">
<article class="task-row">
<header class="task-row-header">
<div class="task-main"></div>
<div class="task-actions"></div>
</header>
<div class="task-progress"></div>
<div class="item-list"></div>
<details class="event-log">
<summary>Activity</summary>
<ol></ol>
</details>
</article>
</template>
<script src="/static/app.js" defer></script>
</body>
</html>
@@ -0,0 +1,131 @@
:root {
color-scheme: light;
font-family: Inter, ui-sans-serif, system-ui, -apple-system, BlinkMacSystemFont, "Segoe UI", sans-serif;
background: #f4f7f8;
color: #18282d;
line-height: 1.4;
}
* { box-sizing: border-box; }
body { margin: 0; min-width: 320px; }
button, input, textarea, select { font: inherit; }
button { cursor: pointer; }
button:disabled { cursor: not-allowed; opacity: .5; }
.topbar {
min-height: 64px;
display: grid;
grid-template-columns: minmax(210px, 1fr) auto minmax(210px, 1fr);
align-items: center;
gap: 20px;
padding: 10px clamp(18px, 4vw, 64px);
background: #10262b;
color: #f9fcfc;
border-bottom: 3px solid #0ac7a1;
}
.brand { display: inline-flex; align-items: center; gap: 10px; font-weight: 750; letter-spacing: 0; }
.brand-mark { display: inline-flex; width: 25px; height: 25px; padding: 4px; gap: 3px; background: #f45e4c; }
.brand-mark i { display: block; width: 4px; background: #10262b; }
.tabs { display: inline-flex; align-self: stretch; }
.tab { border: 0; padding: 0 17px; color: #bad0d1; background: transparent; border-bottom: 2px solid transparent; }
.tab.is-active { color: #ffffff; border-bottom-color: #f5c85d; }
.topbar > .primary-action { justify-self: end; }
main { max-width: 1200px; margin: 0 auto; padding: 38px clamp(18px, 4vw, 56px) 70px; }
.view { display: none; }
.view.is-active { display: block; }
.view-heading { display: flex; justify-content: space-between; align-items: end; gap: 24px; margin-bottom: 25px; }
.eyebrow { margin: 0 0 4px; color: #42747d; font-size: .76rem; font-weight: 750; letter-spacing: .08em; text-transform: uppercase; }
h1, h2, p { margin-top: 0; }
h1 { margin-bottom: 0; font-size: 1.55rem; }
h2 { margin-bottom: 0; font-size: 1.15rem; }
.queue-count { color: #49646b; font-size: .9rem; }
.primary-action, .secondary-action, .icon-button, .close-button {
border: 1px solid transparent;
border-radius: 5px;
min-height: 36px;
padding: 0 14px;
font-weight: 680;
}
.primary-action { background: #0d917c; color: #fff; border-color: #0d917c; }
.primary-action:hover:not(:disabled) { background: #087a69; }
.secondary-action { background: #fff; color: #24434a; border-color: #b9c9cc; }
.secondary-action:hover { background: #edf4f4; }
.icon-button, .close-button { background: transparent; color: #24434a; border-color: #b9c9cc; padding: 0 10px; }
.close-button { color: #ba3629; border-color: #edb9b1; }
.task-list { display: grid; gap: 10px; }
.task-row { background: #fff; border: 1px solid #d4e0e1; border-left: 4px solid #5a7a81; padding: 18px 20px; border-radius: 6px; }
.task-row[data-status="running"] { border-left-color: #0d917c; }
.task-row[data-status="completed"] { border-left-color: #4a8f55; }
.task-row[data-status="partial"], .task-row[data-status="failed"] { border-left-color: #e36b3e; }
.task-row[data-status="cancelled"] { border-left-color: #7d689f; }
.task-row-header { display: flex; justify-content: space-between; gap: 16px; align-items: flex-start; }
.task-main { min-width: 0; }
.task-title { font-weight: 720; overflow-wrap: anywhere; }
.task-meta { margin-top: 3px; color: #5b747a; font-size: .84rem; overflow-wrap: anywhere; }
.task-status { display: inline-block; margin-top: 8px; padding: 2px 8px; background: #e8f1f1; color: #31535a; border-radius: 3px; font-size: .76rem; font-weight: 720; text-transform: uppercase; }
.task-actions { display: inline-flex; flex-wrap: wrap; justify-content: end; gap: 7px; }
.task-progress { margin-top: 14px; }
progress { width: 100%; height: 7px; accent-color: #0d917c; }
.progress-caption { display: flex; justify-content: space-between; margin-bottom: 5px; color: #557076; font-size: .8rem; }
.item-list { display: grid; margin-top: 16px; border-top: 1px solid #e5ecec; }
.item-row { display: grid; grid-template-columns: minmax(0, 1fr) auto; gap: 14px; align-items: center; padding: 11px 0; border-bottom: 1px solid #e5ecec; }
.item-row:last-child { border-bottom: 0; }
.item-name { font-size: .92rem; font-weight: 680; overflow-wrap: anywhere; }
.item-detail { color: #60797e; font-size: .78rem; overflow-wrap: anywhere; }
.item-progress { margin-top: 8px; max-width: 620px; }
.item-progress progress { height: 5px; accent-color: #2378a8; }
.item-warning { color: #a3462b; }
.item-actions { display: flex; align-items: center; gap: 8px; }
.item-actions .secondary-action { min-height: 30px; padding: 0 9px; font-size: .78rem; }
.event-log { margin-top: 14px; color: #587177; font-size: .8rem; }
.event-log summary { cursor: pointer; font-weight: 650; }
.event-log ol { padding-left: 19px; margin-bottom: 0; }
.event-log li { margin-top: 4px; overflow-wrap: anywhere; }
.empty-state { padding: 58px 0; color: #668087; text-align: center; border-top: 1px solid #cddadb; border-bottom: 1px solid #cddadb; }
.settings-form { display: grid; max-width: 720px; gap: 19px; background: #fff; border: 1px solid #d4e0e1; padding: 25px; border-radius: 6px; }
label { display: grid; gap: 7px; color: #345158; font-size: .88rem; font-weight: 700; }
input, textarea, select { width: 100%; border: 1px solid #b9c9cc; border-radius: 4px; background: #fff; color: #18282d; padding: 9px 10px; outline-color: #0d917c; }
textarea { min-height: 220px; resize: vertical; font-family: ui-monospace, SFMono-Regular, Menlo, Consolas, monospace; font-size: .84rem; line-height: 1.45; }
.path-control { display: grid; grid-template-columns: minmax(0, 1fr) auto; gap: 8px; }
.form-actions { display: flex; align-items: center; gap: 13px; }
.form-status, .validation-status { min-height: 1.2em; margin: 0; color: #587177; font-size: .84rem; }
.validation-status.is-invalid { color: #b33c30; }
.validation-status.is-valid { color: #17765d; }
.log-filter { width: 150px; font-size: .78rem; }
.log-filter select { padding: 7px 8px; }
.log-list { display: grid; margin: 0; padding: 0; border-top: 1px solid #cddadb; list-style: none; }
.log-row { display: grid; grid-template-columns: minmax(220px, auto) minmax(0, 1fr); gap: 14px; padding: 10px 3px; border-bottom: 1px solid #dce6e7; font-size: .82rem; }
.log-meta { color: #5b747a; font-family: ui-monospace, SFMono-Regular, Menlo, Consolas, monospace; overflow-wrap: anywhere; }
.log-message { color: #24434a; overflow-wrap: anywhere; }
.log-warning .log-message { color: #9b4b2d; }
.log-error .log-message { color: #b33c30; }
.modal { width: min(740px, calc(100vw - 28px)); border: 0; border-radius: 7px; padding: 0; color: #18282d; box-shadow: 0 22px 80px rgb(17 39 44 / .3); }
.modal::backdrop { background: rgb(10 25 29 / .56); }
.modal form, .directory-modal { padding: 22px; }
.task-modal form { display: grid; gap: 16px; }
.modal-header { display: flex; align-items: flex-start; justify-content: space-between; gap: 16px; }
.modal-footer { display: flex; justify-content: end; gap: 9px; padding-top: 4px; }
.directory-modal { min-height: min(540px, calc(100vh - 42px)); }
.directory-path { margin: 15px 0 10px; font-family: ui-monospace, SFMono-Regular, Menlo, Consolas, monospace; font-size: .82rem; color: #526f75; overflow-wrap: anywhere; }
.directory-actions { display: flex; justify-content: space-between; gap: 10px; padding-bottom: 12px; border-bottom: 1px solid #d6e2e2; }
.directory-list { display: grid; max-height: 360px; overflow: auto; }
.directory-entry { display: flex; width: 100%; min-height: 40px; align-items: center; border: 0; border-bottom: 1px solid #edf1f1; padding: 0 7px; background: #fff; color: #24434a; text-align: left; overflow-wrap: anywhere; }
.directory-entry:hover { background: #eff7f5; }
@media (max-width: 680px) {
.topbar { grid-template-columns: 1fr auto; gap: 10px; padding: 10px 16px; }
.tabs { grid-row: 2; grid-column: 1 / -1; min-height: 36px; }
.topbar > .primary-action { grid-column: 2; grid-row: 1; padding: 0 10px; font-size: .84rem; }
.task-row { padding: 15px; }
.task-row-header { display: grid; }
.task-actions { justify-content: start; }
.item-row { grid-template-columns: minmax(0, 1fr); }
.item-actions { justify-content: flex-start; }
.view-heading { align-items: start; }
.log-row { grid-template-columns: minmax(0, 1fr); gap: 3px; }
}
@@ -0,0 +1,290 @@
"""Thread-safe, process-local task state."""
from __future__ import annotations
import copy
import threading
from datetime import datetime, timezone
from pathlib import Path
from typing import Any, Iterable
from .models import DATE_TEMPLATE, DEFAULT_PROXY, DownloadItemInput
MAX_LOG_ENTRIES = 1_000
LOG_LEVELS = {"debug", "info", "warning", "error"}
def _now() -> str:
return datetime.now(timezone.utc).isoformat()
class TaskState:
"""Keeps settings and task history only for the current process lifetime."""
def __init__(self, download_root: Path, max_active_downloads: int = 2) -> None:
self._lock = threading.RLock()
self._settings = {
"default_title_template": DATE_TEMPLATE,
"default_directory": str(download_root.resolve()),
"proxy": DEFAULT_PROXY,
"max_active_downloads": max_active_downloads,
}
self._tasks: dict[int, dict[str, Any]] = {}
self._logs: list[dict[str, Any]] = []
self._next_task_id = 1
self._next_item_id = 1
self._next_event_id = 1
self._next_log_id = 1
def settings(self) -> dict[str, str | int]:
with self._lock:
return dict(self._settings)
def update_settings(self, values: dict[str, str | int]) -> None:
with self._lock:
self._settings.update(values)
def max_active_downloads(self) -> int:
with self._lock:
return int(self._settings["max_active_downloads"])
def logs(self, level: str | None = None) -> list[dict[str, Any]]:
normalized = (level or "").lower()
if normalized and normalized not in LOG_LEVELS:
raise ValueError("Unknown log level")
with self._lock:
logs = self._logs if not normalized else [entry for entry in self._logs if entry["level"] == normalized]
return copy.deepcopy(logs)
def log(self, level: str, message: str, task_id: int | None = None, item_id: int | None = None) -> None:
with self._lock:
self._log_unlocked(level, message, task_id, item_id)
def create_task(
self,
*,
title: str,
base_dir: Path,
output_dir: Path,
proxy: str,
items: Iterable[DownloadItemInput],
) -> int:
with self._lock:
task_id = self._next_task_id
self._next_task_id += 1
task = {
"id": task_id,
"title": title,
"base_dir": str(base_dir),
"output_dir": str(output_dir),
"proxy": proxy,
"status": "queued",
"cancel_requested": False,
"created_at": _now(),
"started_at": None,
"completed_at": None,
"error": None,
"items": [],
"events": [],
}
for position, item in enumerate(items, 1):
item_id = self._next_item_id
self._next_item_id += 1
task["items"].append(
{
"id": item_id,
"task_id": task_id,
"position": position,
"title": item.title,
"image_src": item.image_src,
"m3u8_url": item.m3u8_url,
"m3u8_referer": item.m3u8_referer,
"status": "queued",
"stage": "Queued",
"completed_segments": 0,
"total_segments": 0,
"message": "",
"error": None,
"warning": None,
"output_path": None,
"cover_path": None,
"started_at": None,
"completed_at": None,
}
)
self._event_unlocked(task, None, "info", "Task queued")
self._tasks[task_id] = task
return task_id
def claim_next_item(self) -> tuple[dict[str, Any], dict[str, Any]] | None:
with self._lock:
for task in self._tasks.values():
if task["status"] == "queued":
task.update(status="running", started_at=_now(), cancel_requested=False)
self._event_unlocked(task, None, "info", "Task started")
if task["status"] != "running" or task["cancel_requested"]:
continue
for item in task["items"]:
if item["status"] != "queued":
continue
item.update(status="running", stage="Preparing", started_at=_now(), message="Preparing download")
self._event_unlocked(task, int(item["id"]), "info", "Video started")
return copy.deepcopy(task), copy.deepcopy(item)
return None
def get_task(self, task_id: int) -> dict[str, Any] | None:
with self._lock:
task = self._tasks.get(task_id)
return copy.deepcopy(task) if task else None
def list_tasks(self) -> list[dict[str, Any]]:
with self._lock:
return copy.deepcopy(list(reversed(list(self._tasks.values()))))
def update_item(self, item_id: int, **values: Any) -> None:
with self._lock:
item = self._find_item_unlocked(item_id)
item.update(values)
def event(self, task_id: int, item_id: int | None, level: str, message: str) -> None:
with self._lock:
self._event_unlocked(self._tasks[task_id], item_id, level, message)
def cancel_task(self, task_id: int) -> bool:
with self._lock:
task = self._tasks.get(task_id)
if task is None or task["status"] in {"completed", "partial", "failed", "cancelled"}:
return False
now = _now()
task["cancel_requested"] = True
for item in task["items"]:
if item["status"] == "queued":
item.update(status="cancelled", stage="Cancelled", completed_at=now)
self._event_unlocked(task, None, "info", "Cancellation requested")
self._finalize_task_unlocked(task)
return True
def cancel_requested(self, task_id: int) -> bool:
with self._lock:
return bool(self._tasks[task_id]["cancel_requested"])
def finalize_task_if_ready(self, task_id: int) -> None:
with self._lock:
task = self._tasks[task_id]
self._finalize_task_unlocked(task)
def retry_task(self, task_id: int) -> bool:
with self._lock:
task = self._tasks.get(task_id)
if task is None or task["status"] == "running":
return False
retryable = [item for item in task["items"] if item["status"] in {"failed", "cancelled", "completed_warning"}]
if not retryable:
return False
for item in retryable:
item.update(
status="queued",
stage="Queued",
completed_segments=0,
total_segments=0,
message="",
error=None,
warning=None,
started_at=None,
completed_at=None,
)
task.update(status="queued", cancel_requested=False, started_at=None, completed_at=None, error=None)
self._event_unlocked(task, None, "info", "Task retry queued")
return True
def retry_item(self, item_id: int) -> bool:
with self._lock:
task, item = self._find_task_and_item_unlocked(item_id)
if task["status"] == "running" or item["status"] == "completed":
return False
item.update(
status="queued",
stage="Cover retry" if item["status"] == "completed_warning" else "Queued",
completed_segments=0,
total_segments=0,
message="",
error=None,
warning=None,
started_at=None,
completed_at=None,
)
task.update(status="queued", cancel_requested=False, started_at=None, completed_at=None, error=None)
self._event_unlocked(task, item_id, "info", "Video retry queued")
return True
def _find_item_unlocked(self, item_id: int) -> dict[str, Any]:
return self._find_task_and_item_unlocked(item_id)[1]
def _find_task_and_item_unlocked(self, item_id: int) -> tuple[dict[str, Any], dict[str, Any]]:
for task in self._tasks.values():
for item in task["items"]:
if item["id"] == item_id:
return task, item
raise KeyError(item_id)
def _event_unlocked(self, task: dict[str, Any], item_id: int | None, level: str, message: str) -> None:
created_at = _now()
task["events"].insert(
0,
{
"id": self._next_event_id,
"task_id": task["id"],
"item_id": item_id,
"created_at": created_at,
"level": level,
"message": message[:2000],
},
)
del task["events"][30:]
self._next_event_id += 1
self._log_unlocked(level, message, int(task["id"]), item_id, created_at)
def _log_unlocked(
self,
level: str,
message: str,
task_id: int | None = None,
item_id: int | None = None,
created_at: str | None = None,
) -> None:
normalized = level.lower()
if normalized not in LOG_LEVELS:
normalized = "info"
self._logs.insert(
0,
{
"id": self._next_log_id,
"created_at": created_at or _now(),
"level": normalized,
"task_id": task_id,
"item_id": item_id,
"message": message[:2000],
},
)
del self._logs[MAX_LOG_ENTRIES:]
self._next_log_id += 1
def _finalize_task_unlocked(self, task: dict[str, Any]) -> None:
if task["status"] not in {"queued", "running"}:
return
statuses = {item["status"] for item in task["items"]}
if statuses & {"queued", "running"}:
return
if task["cancel_requested"] and statuses <= {"completed", "completed_warning", "cancelled"}:
status = "cancelled"
elif statuses <= {"completed", "completed_warning"}:
status = "completed"
elif statuses & {"completed", "completed_warning"}:
status = "partial"
elif statuses == {"cancelled"}:
status = "cancelled"
else:
status = "failed"
task.update(status=status, completed_at=_now())
self._event_unlocked(task, None, "info" if status == "completed" else "warning", f"Task {status}")
+340
View File
@@ -0,0 +1,340 @@
"""Small dependency-free HTTP server and JSON API for the downloader UI."""
from __future__ import annotations
import argparse
import json
import mimetypes
from dataclasses import dataclass
from http import HTTPStatus
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
from pathlib import Path
from typing import Any
from urllib.parse import parse_qs, unquote, urlsplit
from .models import (
MAX_ACTIVE_DOWNLOADS,
ValidationError,
decode_download_items,
ensure_within_root,
resolve_default_title,
validate_folder_title,
validate_max_active_downloads,
validate_proxy,
)
from .tasks import TaskState
from .worker import TaskRunner
MAX_REQUEST_BYTES = 2 * 1024 * 1024
@dataclass(frozen=True)
class ServiceConfig:
host: str
port: int
download_root: Path
max_active_downloads: int
class DownloadService:
def __init__(self, config: ServiceConfig) -> None:
root = config.download_root.expanduser().resolve()
if not root.is_dir():
raise RuntimeError(f"Download root does not exist or is not a directory: {root}")
self.config = ServiceConfig(config.host, config.port, root, config.max_active_downloads)
self.store = TaskState(root, self.config.max_active_downloads)
self.runner = TaskRunner(self.store)
def start(self) -> None:
self.runner.start()
def stop(self) -> None:
self.runner.stop()
def bootstrap(self) -> dict[str, Any]:
settings = self.store.settings()
return {
"settings": settings,
"default_title": resolve_default_title(str(settings["default_title_template"])),
"download_root": str(self.config.download_root),
"tasks": [_public_task(task) for task in self.store.list_tasks()],
}
def update_settings(self, value: dict[str, Any]) -> dict[str, str | int]:
current = self.store.settings()
template = value.get("default_title_template", current["default_title_template"])
directory = value.get("default_directory", current["default_directory"])
proxy = value.get("proxy", current["proxy"])
max_active_downloads = value.get("max_active_downloads", current["max_active_downloads"])
if not isinstance(template, str):
raise ValidationError("Default title must be a string")
selected = ensure_within_root(directory, self.config.download_root)
normalized_proxy = validate_proxy(proxy)
normalized_max_active_downloads = validate_max_active_downloads(max_active_downloads)
settings = {
"default_title_template": template,
"default_directory": str(selected),
"proxy": normalized_proxy,
"max_active_downloads": normalized_max_active_downloads,
}
self.store.update_settings(settings)
self.store.log("info", f"Settings updated: concurrent downloads={normalized_max_active_downloads}")
return settings
def logs(self, level: str | None) -> dict[str, Any]:
try:
return {"logs": self.store.logs(level)}
except ValueError as error:
raise ValidationError(str(error)) from error
def directories(self, requested: str | None) -> dict[str, Any]:
selected = ensure_within_root(requested or self.config.download_root, self.config.download_root)
directories: list[dict[str, str]] = []
for entry in sorted(selected.iterdir(), key=lambda path: path.name.casefold()):
try:
resolved = entry.resolve()
resolved.relative_to(self.config.download_root)
except (OSError, ValueError):
continue
if resolved.is_dir():
directories.append({"name": entry.name, "path": str(resolved)})
parent: str | None = None
if selected != self.config.download_root:
candidate = selected.parent.resolve()
try:
candidate.relative_to(self.config.download_root)
except ValueError:
pass
else:
parent = str(candidate)
return {"path": str(selected), "parent": parent, "directories": directories}
def validate_payload(self, payload: object) -> dict[str, Any]:
items = decode_download_items(payload if isinstance(payload, str) else "")
return {"valid": True, "count": len(items)}
def create_task(self, value: dict[str, Any]) -> dict[str, Any]:
payload = value.get("payload")
title = validate_folder_title(value.get("title", ""))
items = decode_download_items(payload if isinstance(payload, str) else "")
settings = self.store.settings()
base_dir = ensure_within_root(value.get("directory", settings["default_directory"]), self.config.download_root)
output_dir = base_dir if not title else ensure_within_root(base_dir / title, self.config.download_root, must_exist=False)
output_dir.mkdir(parents=True, exist_ok=True)
task_id = self.store.create_task(
title=title,
base_dir=base_dir,
output_dir=output_dir,
proxy=settings["proxy"],
items=items,
)
task = self.store.get_task(task_id)
assert task is not None
return _public_task(task)
def get_task(self, task_id: int) -> dict[str, Any] | None:
task = self.store.get_task(task_id)
return _public_task(task) if task else None
def list_tasks(self) -> list[dict[str, Any]]:
return [_public_task(task) for task in self.store.list_tasks()]
class ServiceRequestHandler(BaseHTTPRequestHandler):
server: "DownloadHTTPServer"
def do_GET(self) -> None: # noqa: N802
parsed = urlsplit(self.path)
if parsed.path == "/":
self._serve_static("index.html")
return
if parsed.path.startswith("/static/"):
self._serve_static(unquote(parsed.path.removeprefix("/static/")))
return
if parsed.path == "/api/bootstrap":
self._json(HTTPStatus.OK, self.server.service.bootstrap())
return
if parsed.path == "/api/logs":
query = parse_qs(parsed.query)
self._call(lambda: self.server.service.logs(query.get("level", [None])[0]))
return
if parsed.path == "/api/tasks":
self._json(HTTPStatus.OK, {"tasks": self.server.service.list_tasks()})
return
if parsed.path.startswith("/api/tasks/"):
task_id = _id_from_path(parsed.path, "/api/tasks/")
task = self.server.service.get_task(task_id) if task_id is not None else None
if task is None:
self._error(HTTPStatus.NOT_FOUND, "Task not found")
else:
self._json(HTTPStatus.OK, task)
return
if parsed.path == "/api/directories":
query = parse_qs(parsed.query)
self._call(lambda: self.server.service.directories(query.get("path", [None])[0]))
return
self._error(HTTPStatus.NOT_FOUND, "Not found")
def do_POST(self) -> None: # noqa: N802
value = self._request_json()
if value is None:
return
if self.path == "/api/validate":
self._call(lambda: self.server.service.validate_payload(value.get("payload")))
return
if self.path == "/api/tasks":
self._call(lambda: self.server.service.create_task(value), status=HTTPStatus.CREATED)
return
if self.path.endswith("/cancel") and self.path.startswith("/api/tasks/"):
task_id = _id_from_path(self.path[: -len("/cancel")], "/api/tasks/")
if task_id is None or not self.server.service.store.cancel_task(task_id):
self._error(HTTPStatus.CONFLICT, "Task cannot be cancelled")
else:
task = self.server.service.get_task(task_id)
assert task is not None
self._json(HTTPStatus.OK, task)
return
if self.path.endswith("/retry") and self.path.startswith("/api/tasks/"):
task_id = _id_from_path(self.path[: -len("/retry")], "/api/tasks/")
if task_id is None or not self.server.service.store.retry_task(task_id):
self._error(HTTPStatus.CONFLICT, "Task cannot be retried")
else:
task = self.server.service.get_task(task_id)
assert task is not None
self._json(HTTPStatus.OK, task)
return
if self.path.endswith("/retry") and self.path.startswith("/api/items/"):
item_id = _id_from_path(self.path[: -len("/retry")], "/api/items/")
if item_id is None or not self.server.service.store.retry_item(item_id):
self._error(HTTPStatus.CONFLICT, "Video cannot be retried")
else:
self._json(HTTPStatus.OK, {"ok": True})
return
self._error(HTTPStatus.NOT_FOUND, "Not found")
def do_PUT(self) -> None: # noqa: N802
if self.path != "/api/settings":
self._error(HTTPStatus.NOT_FOUND, "Not found")
return
value = self._request_json()
if value is not None:
self._call(lambda: self.server.service.update_settings(value))
def _call(self, callback: Any, *, status: HTTPStatus = HTTPStatus.OK) -> None:
try:
self._json(status, callback())
except ValidationError as error:
self._error(HTTPStatus.BAD_REQUEST, str(error))
except OSError as error:
self._error(HTTPStatus.BAD_REQUEST, str(error))
def _request_json(self) -> dict[str, Any] | None:
try:
length = int(self.headers.get("Content-Length", "0"))
except ValueError:
self._error(HTTPStatus.BAD_REQUEST, "Invalid Content-Length")
return None
if length <= 0 or length > MAX_REQUEST_BYTES:
self._error(HTTPStatus.REQUEST_ENTITY_TOO_LARGE, "Request body must be between 1 byte and 2 MiB")
return None
try:
value = json.loads(self.rfile.read(length).decode("utf-8"))
except (UnicodeDecodeError, json.JSONDecodeError):
self._error(HTTPStatus.BAD_REQUEST, "Request body must be UTF-8 JSON")
return None
if not isinstance(value, dict):
self._error(HTTPStatus.BAD_REQUEST, "Request body must be a JSON object")
return None
return value
def _serve_static(self, relative: str) -> None:
static_root = Path(__file__).with_name("static").resolve()
target = (static_root / relative).resolve()
try:
target.relative_to(static_root)
except ValueError:
self._error(HTTPStatus.NOT_FOUND, "Not found")
return
if not target.is_file():
self._error(HTTPStatus.NOT_FOUND, "Not found")
return
content_type = mimetypes.guess_type(target.name)[0] or "application/octet-stream"
data = target.read_bytes()
self.send_response(HTTPStatus.OK)
self.send_header("Content-Type", f"{content_type}; charset=utf-8" if content_type.startswith("text/") else content_type)
self.send_header("Content-Length", str(len(data)))
self.send_header("Cache-Control", "no-cache")
self.end_headers()
self.wfile.write(data)
def _json(self, status: HTTPStatus, value: object) -> None:
data = json.dumps(value, ensure_ascii=False, default=str).encode("utf-8")
self.send_response(status)
self.send_header("Content-Type", "application/json; charset=utf-8")
self.send_header("Content-Length", str(len(data)))
self.send_header("Cache-Control", "no-store")
self.end_headers()
self.wfile.write(data)
def _error(self, status: HTTPStatus, message: str) -> None:
self._json(status, {"error": message})
def log_message(self, format: str, *args: object) -> None:
if not self.path.startswith("/api/bootstrap"):
self.server.service.store.log("debug", f"HTTP {self.command} {self.path}: {format % args}")
class DownloadHTTPServer(ThreadingHTTPServer):
def __init__(self, config: ServiceConfig, service: DownloadService) -> None:
super().__init__((config.host, config.port), ServiceRequestHandler)
self.service = service
def serve(config: ServiceConfig) -> None:
service = DownloadService(config)
server = DownloadHTTPServer(config, service)
service.start()
try:
server.serve_forever(poll_interval=0.25)
except KeyboardInterrupt:
pass
finally:
server.shutdown()
service.stop()
server.server_close()
def build_argument_parser() -> argparse.ArgumentParser:
parser = argparse.ArgumentParser(description="Web service for N_m3u8DL-PY")
parser.add_argument("--host", default="0.0.0.0")
parser.add_argument("--port", type=int, default=8000)
parser.add_argument("--download-root", type=Path, default=Path("/tmp"))
parser.add_argument("--max-active-downloads", type=int, default=2)
return parser
def config_from_args(args: argparse.Namespace) -> ServiceConfig:
if not 1 <= args.port <= 65535:
raise ValueError("--port must be between 1 and 65535")
if not 1 <= args.max_active_downloads <= MAX_ACTIVE_DOWNLOADS:
raise ValueError(f"--max-active-downloads must be between 1 and {MAX_ACTIVE_DOWNLOADS}")
return ServiceConfig(args.host, args.port, args.download_root, args.max_active_downloads)
def _id_from_path(path: str, prefix: str) -> int | None:
value = path.removeprefix(prefix)
try:
return int(value)
except ValueError:
return None
def _public_task(task: dict[str, Any] | None) -> dict[str, Any]:
if task is None:
return {}
public = dict(task)
public["items"] = [dict(item) for item in public.get("items", [])]
for item in public["items"]:
item.pop("m3u8_referer", None)
return public
@@ -0,0 +1,417 @@
"""Durable task scheduler and media post-processing pipeline."""
from __future__ import annotations
import concurrent.futures
import os
import re
import shutil
import subprocess
import tempfile
import threading
import time
from pathlib import Path
from urllib.error import HTTPError, URLError
from urllib.parse import urlsplit, urlunsplit
from urllib.request import Request
from n_m3u8dl_py import DownloadCancelledError, DownloadEventKind, DownloadRequest, N_m3u8DL, RequestOptions
from n_m3u8dl_py.http import build_proxy_opener
from n_m3u8dl_py.utils import valid_filename
from .models import MAX_ACTIVE_DOWNLOADS
from .tasks import TaskState
class TaskRunner:
"""Runs queued video downloads while keeping HTTP request handling independent."""
def __init__(self, store: TaskState, worker_capacity: int = MAX_ACTIVE_DOWNLOADS, ffmpeg: str | None = None) -> None:
self.store = store
self.worker_capacity = max(1, worker_capacity)
self.ffmpeg = ffmpeg or shutil.which("ffmpeg")
if not self.ffmpeg:
raise RuntimeError("ffmpeg is required to run m3u8downloaderd")
self._stop = threading.Event()
self._dispatcher: threading.Thread | None = None
self._executor = concurrent.futures.ThreadPoolExecutor(
max_workers=self.worker_capacity,
thread_name_prefix="m3u8downloaderd-download",
)
def start(self) -> None:
if self._dispatcher is None:
self._dispatcher = threading.Thread(target=self._dispatch, name="m3u8downloaderd-dispatch", daemon=True)
self._dispatcher.start()
def stop(self) -> None:
self._stop.set()
if self._dispatcher is not None:
self._dispatcher.join(timeout=5)
self._executor.shutdown(wait=False, cancel_futures=True)
def _dispatch(self) -> None:
active: set[concurrent.futures.Future[None]] = set()
while not self._stop.is_set():
completed = {future for future in active if future.done()}
active.difference_update(completed)
for future in completed:
try:
future.result()
except Exception as error:
self.store.log("error", f"Download worker crashed: {error}")
while len(active) < self.store.max_active_downloads():
claimed = self.store.claim_next_item()
if claimed is None:
break
task, item = claimed
active.add(self._executor.submit(self._run_item, task, item))
self._stop.wait(0.2)
def _run_item(self, task: dict[str, object], item: dict[str, object]) -> None:
item_id = int(item["id"])
task_id = int(task["id"])
try:
if item["stage"] == "Cover retry" and item["output_path"]:
self._retry_cover(task, item)
return
self._download_and_finalize(task, item)
except DownloadCancelledError:
self.store.update_item(
item_id,
status="cancelled",
stage="Cancelled",
completed_at=_timestamp(),
message="Cancellation completed",
)
self.store.event(task_id, item_id, "info", "Video cancelled")
except Exception as error:
self.store.update_item(
item_id,
status="failed",
stage="Failed",
error=str(error),
completed_at=_timestamp(),
message="Download failed",
)
self.store.event(task_id, item_id, "error", str(error))
finally:
self.store.finalize_task_if_ready(task_id)
def _download_and_finalize(self, task: dict[str, object], item: dict[str, object]) -> None:
task_id = int(task["id"])
item_id = int(item["id"])
output_dir = Path(str(task["output_dir"]))
output_dir.mkdir(parents=True, exist_ok=True)
progress: dict[str, tuple[int, int]] = {}
def on_event(event: object) -> bool | None:
if self.store.cancel_requested(task_id):
return False
kind = getattr(event, "kind", None)
track_id = getattr(event, "track_id", None) or "unknown"
if kind is DownloadEventKind.TRACK_STARTED:
progress[track_id] = (0, int(getattr(event, "total_segments", 0) or 0))
elif kind is DownloadEventKind.SEGMENT_COMPLETED:
progress[track_id] = (
int(getattr(event, "completed_segments", 0) or 0),
int(getattr(event, "total_segments", 0) or 0),
)
elif kind is DownloadEventKind.TRACK_COMPLETED and track_id in progress:
done, total = progress[track_id]
progress[track_id] = (total or done, total or done)
if progress:
completed = sum(value[0] for value in progress.values())
total = sum(value[1] for value in progress.values())
self.store.update_item(
item_id,
stage="Downloading",
completed_segments=completed,
total_segments=total,
message=getattr(event, "message", "Downloading"),
)
message = str(getattr(event, "message", ""))
if message:
level = "debug" if kind is DownloadEventKind.SEGMENT_COMPLETED else "info"
self.store.log(level, message, task_id, item_id)
return None
self.store.log("info", "Inspecting media playlist", task_id, item_id)
self.store.update_item(item_id, stage="Downloading", message="Inspecting media playlist")
options = RequestOptions(
headers={"Referer": str(item["m3u8_referer"])},
proxy=str(task["proxy"]) or None,
use_system_proxy=False,
)
with tempfile.TemporaryDirectory(prefix=f".m3u8downloaderd-{task_id}-{item_id}-", dir=output_dir) as temporary_dir:
request = DownloadRequest(
output_dir=output_dir,
temporary_dir=temporary_dir,
file_name=str(item["title"]),
save_pattern=f"<SaveName>.source_{item_id}_<MediaType>.<Ext>",
thread_count=max(1, os.cpu_count() or 1),
ffmpeg_path=self.ffmpeg,
)
result = N_m3u8DL(options).download_url(str(item["m3u8_url"]), request, on_event=on_event)
if self.store.cancel_requested(task_id):
raise DownloadCancelledError("Download cancelled by user")
self.store.log("info", "Creating MP4", task_id, item_id)
self.store.update_item(item_id, stage="Muxing", message="Creating MP4")
output_path = self._reserve_output_path(output_dir, str(item["title"]))
mux_temp = output_dir / f".{output_path.stem}.{item_id}.mux.tmp.mp4"
try:
self._mux_mp4([file.path for file in result.files], mux_temp)
for source in result.files:
source.path.unlink(missing_ok=True)
self._complete_cover_step(task, item, output_path, mux_temp)
except Exception:
output_path.unlink(missing_ok=True)
mux_temp.unlink(missing_ok=True)
raise
def _retry_cover(self, task: dict[str, object], item: dict[str, object]) -> None:
item_id = int(item["id"])
output_path = Path(str(item["output_path"]))
if not output_path.is_file():
raise FileNotFoundError("The existing MP4 for this cover retry no longer exists")
self.store.log("info", "Retrying cover", int(task["id"]), item_id)
self.store.update_item(item_id, stage="Cover", message="Retrying cover")
cover_path = Path(str(item["cover_path"])) if item["cover_path"] else None
if cover_path is None or not cover_path.is_file():
cover_path = self._download_cover(
str(item["image_src"]),
output_path.parent,
item_id,
str(task["proxy"]),
task_id=int(task["id"]),
)
covered_temp = output_path.parent / f".{output_path.stem}.{item_id}.cover.tmp.mp4"
try:
self._attach_cover(output_path, cover_path, covered_temp)
os.replace(covered_temp, output_path)
cover_path.unlink(missing_ok=True)
self.store.update_item(
item_id,
status="completed",
stage="Completed",
warning=None,
cover_path=None,
completed_at=_timestamp(),
message="Video and cover completed",
)
except Exception as error:
retained_cover = self._retain_cover(cover_path, output_path)
self.store.update_item(
item_id,
status="completed_warning",
stage="Completed with warning",
warning=str(error),
cover_path=str(retained_cover) if retained_cover else None,
completed_at=_timestamp(),
message="Video completed without an embedded cover",
)
self.store.event(int(task["id"]), item_id, "warning", f"Cover retry failed: {error}")
finally:
covered_temp.unlink(missing_ok=True)
def _complete_cover_step(
self,
task: dict[str, object],
item: dict[str, object],
output_path: Path,
mux_temp: Path,
) -> None:
item_id = int(item["id"])
self.store.log("info", "Downloading cover image", int(task["id"]), item_id)
self.store.update_item(item_id, stage="Cover", message="Downloading cover image")
cover_path: Path | None = None
covered_temp = output_path.parent / f".{output_path.stem}.{item_id}.cover.tmp.mp4"
try:
cover_path = self._download_cover(
str(item["image_src"]),
output_path.parent,
item_id,
str(task["proxy"]),
task_id=int(task["id"]),
)
self.store.log("info", "Embedding cover image", int(task["id"]), item_id)
self.store.update_item(item_id, stage="Cover", message="Embedding cover image")
self._attach_cover(mux_temp, cover_path, covered_temp)
os.replace(covered_temp, output_path)
mux_temp.unlink(missing_ok=True)
cover_path.unlink(missing_ok=True)
self.store.update_item(
item_id,
status="completed",
stage="Completed",
output_path=str(output_path),
cover_path=None,
completed_at=_timestamp(),
message="Video and cover completed",
)
self.store.event(int(task["id"]), item_id, "info", f"Saved {output_path.name}")
except Exception as error:
os.replace(mux_temp, output_path)
retained_cover = self._retain_cover(cover_path, output_path)
self.store.update_item(
item_id,
status="completed_warning",
stage="Completed with warning",
output_path=str(output_path),
cover_path=str(retained_cover) if retained_cover else None,
warning=str(error),
completed_at=_timestamp(),
message="Video completed without an embedded cover",
)
self.store.event(int(task["id"]), item_id, "warning", f"Cover step failed: {error}")
finally:
covered_temp.unlink(missing_ok=True)
def _mux_mp4(self, source_paths: list[Path], target: Path) -> None:
sources = [path for path in source_paths if path.is_file()]
if not sources:
raise RuntimeError("Downloader did not produce media files")
command = [self.ffmpeg, "-hide_banner", "-loglevel", "error", "-nostdin", "-y"]
for source in sources:
command.extend(["-i", str(source)])
for index in range(len(sources)):
command.extend(["-map", str(index)])
command.extend(["-c", "copy"])
if any(path.suffix.lower() == ".srt" for path in sources):
command.extend(["-c:s", "mov_text"])
command.extend(["-movflags", "+faststart", str(target)])
self._run_ffmpeg(command, "MP4 muxing failed")
def _attach_cover(self, video: Path, cover: Path, target: Path) -> None:
command = [
self.ffmpeg,
"-hide_banner",
"-loglevel",
"error",
"-nostdin",
"-y",
"-i",
str(video),
"-i",
str(cover),
"-map",
"0",
"-map",
"1:v:0",
"-c",
"copy",
"-c:v:1",
"mjpeg",
"-disposition:v:1",
"attached_pic",
"-movflags",
"+faststart",
str(target),
]
self._run_ffmpeg(command, "Cover embedding failed")
def _download_cover(self, url: str, output_dir: Path, item_id: int, proxy: str, *, task_id: int) -> Path:
preferred_url = _preferred_cover_url(url)
if preferred_url != url:
try:
return self._download_cover_url(preferred_url, output_dir, item_id, proxy)
except Exception as error:
self.store.log(
"warning",
f"Preferred s1080 cover URL download failed; falling back to the original URL: {error}",
task_id,
item_id,
)
return self._download_cover_url(url, output_dir, item_id, proxy)
def _download_cover_url(self, url: str, output_dir: Path, item_id: int, proxy: str) -> Path:
request = Request(url, headers={"User-Agent": "m3u8downloaderd/0.1"}, method="GET")
opener = build_proxy_opener(proxy or None, use_system_proxy=False)
for attempt in range(4):
try:
with opener.open(request, timeout=100) as response:
content_type = response.headers.get_content_type()
data = _read_limited(response, 25 * 1024 * 1024)
break
except (HTTPError, URLError, TimeoutError, OSError):
if attempt == 3:
raise
time.sleep(attempt + 1)
extension = _image_extension(data, content_type, url)
path = output_dir / f".m3u8downloaderd-{item_id}.cover{extension}"
path.write_bytes(data)
return path
@staticmethod
def _retain_cover(cover_path: Path | None, output_path: Path) -> Path | None:
if cover_path is None or not cover_path.is_file():
return None
retained = output_path.with_name(f"{output_path.stem}.cover{cover_path.suffix}")
if retained.exists():
retained = output_path.with_name(f"{output_path.stem}.cover_{int(time.time())}{cover_path.suffix}")
cover_path.replace(retained)
return retained
@staticmethod
def _reserve_output_path(output_dir: Path, title: str) -> Path:
stem = valid_filename(title, 180)
suffix = 0
while True:
candidate = output_dir / f"{stem}{'' if suffix == 0 else f'_{suffix}'}.mp4"
try:
descriptor = os.open(candidate, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o644)
except FileExistsError:
suffix += 1
continue
os.close(descriptor)
return candidate
@staticmethod
def _run_ffmpeg(command: list[str], summary: str) -> None:
result = subprocess.run(command, capture_output=True, text=True)
if result.returncode:
detail = result.stderr.strip()[-1500:] or summary
raise RuntimeError(f"{summary}: {detail}")
def _read_limited(response: object, limit: int) -> bytes:
chunks: list[bytes] = []
remaining = limit + 1
while remaining:
chunk = response.read(min(1024 * 1024, remaining)) # type: ignore[attr-defined]
if not chunk:
break
chunks.append(chunk)
remaining -= len(chunk)
data = b"".join(chunks)
if len(data) > limit:
raise ValueError("Cover image exceeds the 25 MiB size limit")
return data
def _image_extension(data: bytes, content_type: str, url: str) -> str:
if data.startswith(b"\xff\xd8\xff") or content_type == "image/jpeg":
return ".jpg"
if data.startswith(b"\x89PNG\r\n\x1a\n") or content_type == "image/png":
return ".png"
if data.startswith(b"RIFF") and data[8:12] == b"WEBP":
return ".webp"
suffix = Path(urlsplit(url).path).suffix.lower()
return suffix if suffix in {".gif", ".bmp", ".avif"} else ".img"
_IMG2_SIZE_PATH = re.compile(r"/img2/s\d+/")
def _preferred_cover_url(url: str) -> str:
parts = urlsplit(url)
preferred_path = _IMG2_SIZE_PATH.sub("/img2/s1080/", parts.path)
if preferred_path == parts.path:
return url
return urlunsplit(parts._replace(path=preferred_path))
def _timestamp() -> str:
from datetime import datetime, timezone
return datetime.now(timezone.utc).isoformat()
+8
View File
@@ -0,0 +1,8 @@
from __future__ import annotations
import sys
from pathlib import Path
ROOT = Path(__file__).parents[2]
sys.path[:0] = [str(ROOT / "m3u8downloaderd" / "src"), str(ROOT / "N_m3u8DL_py" / "src")]
+62
View File
@@ -0,0 +1,62 @@
from __future__ import annotations
import base64
import json
from datetime import date, datetime
import pytest
from m3u8downloaderd.models import DATE_TEMPLATE, DEFAULT_PROXY, MAX_ACTIVE_DOWNLOADS, ValidationError, decode_download_items, ensure_within_root, resolve_default_title, validate_max_active_downloads, validate_proxy
def _payload() -> list[dict[str, object]]:
return [
{
"code": 17,
"title": "Episode 1",
"href": "/episode/1",
"image_src": "https://cdn.example.test/cover.jpg",
"m3u8_url": "https://cdn.example.test/video.m3u8",
"m3u8_referer": "https://example.test/watch/1",
"ignored": True,
}
]
def test_download_items_accept_json_and_url_safe_base64() -> None:
value = json.dumps(_payload())
encoded = base64.urlsafe_b64encode(value.encode()).decode().rstrip("=")
for input_value in (value, encoded):
items = decode_download_items(input_value)
assert len(items) == 1
assert items[0].title == "Episode 1"
assert items[0].code == 17
def test_download_items_reject_non_http_sources() -> None:
payload = _payload()
payload[0]["m3u8_url"] = "file:///etc/passwd"
with pytest.raises(ValidationError, match="HTTP"):
decode_download_items(json.dumps(payload))
def test_template_proxy_and_root_rules(tmp_path) -> None:
assert resolve_default_title(DATE_TEMPLATE, datetime(2026, 9, 26, 14, 30, 5)) == "2026-09-26_14-30-05"
assert resolve_default_title("YYYY-MM-DD", date(2026, 9, 26)) == "2026-09-26"
assert resolve_default_title("weekly") == "weekly"
assert validate_proxy("") == ""
assert validate_proxy("http://127.0.0.1:7890") == "http://127.0.0.1:7890"
assert validate_proxy(DEFAULT_PROXY) == DEFAULT_PROXY
assert validate_max_active_downloads("2") == 2
assert validate_max_active_downloads(MAX_ACTIVE_DOWNLOADS) == MAX_ACTIVE_DOWNLOADS
with pytest.raises(ValidationError, match="Authenticated"):
validate_proxy("http://name:secret@127.0.0.1:7890")
with pytest.raises(ValidationError, match="between"):
validate_max_active_downloads(MAX_ACTIVE_DOWNLOADS + 1)
child = tmp_path / "child"
child.mkdir()
assert ensure_within_root(child, tmp_path) == child.resolve()
with pytest.raises(ValidationError, match="inside"):
ensure_within_root(tmp_path.parent, tmp_path)
+307
View File
@@ -0,0 +1,307 @@
from __future__ import annotations
import contextlib
import functools
import http.server
import json
import subprocess
import threading
import time
from pathlib import Path
from urllib.request import Request, urlopen
from m3u8downloaderd.models import DATE_TEMPLATE, DEFAULT_PROXY, decode_download_items
from m3u8downloaderd.tasks import TaskState
from m3u8downloaderd.web import DownloadHTTPServer, DownloadService, ServiceConfig, build_argument_parser, config_from_args
from m3u8downloaderd.worker import TaskRunner
def _item(base_url: str) -> dict[str, str]:
return {
"code": "episode-1",
"title": "episode",
"href": "/episode-1",
"image_src": f"{base_url}/cover.jpg",
"m3u8_url": f"{base_url}/video.m3u8",
"m3u8_referer": "https://example.test/watch/episode-1",
}
def test_http_api_validates_and_snapshots_settings(tmp_path: Path) -> None:
config = ServiceConfig("127.0.0.1", 0, tmp_path / "downloads", 1)
config.download_root.mkdir()
service = DownloadService(config)
assert service.bootstrap()["settings"]["proxy"] == DEFAULT_PROXY
assert service.bootstrap()["settings"]["max_active_downloads"] == 1
server = DownloadHTTPServer(config, service)
thread = threading.Thread(target=server.serve_forever, daemon=True)
thread.start()
try:
base = f"http://127.0.0.1:{server.server_port}"
invalid = _request(base, "/api/validate", {"payload": "[]"})
assert invalid[0] == 400
assert "non-empty" in invalid[1]["error"]
validation = _request(base, "/api/validate", {"payload": json.dumps([_item("http://127.0.0.1")])})
assert validation == (200, {"valid": True, "count": 1})
settings = _request(
base,
"/api/settings",
{
"default_title_template": "YYYY-MM-DD",
"default_directory": str(config.download_root),
"proxy": "http://127.0.0.1:7890",
"max_active_downloads": 2,
},
method="PUT",
)
assert settings[0] == 200
assert settings[1]["max_active_downloads"] == 2
assert service.store.max_active_downloads() == 2
task = _request(base, "/api/tasks", {"title": "batch", "directory": str(config.download_root), "payload": json.dumps([_item("http://127.0.0.1")])})
assert task[0] == 201
assert task[1]["proxy"] == "http://127.0.0.1:7890"
assert task[1]["output_dir"] == str(config.download_root / "batch")
assert "m3u8_referer" not in task[1]["items"][0]
service.store.log("error", "Expected log entry")
logs = _get(base, "/api/logs?level=error")
assert logs[0] == 200
assert logs[1]["logs"][0]["message"] == "Expected log entry"
finally:
server.shutdown()
server.server_close()
thread.join()
service.stop()
def test_runner_creates_mp4_with_attached_cover(tmp_path: Path) -> None:
source = tmp_path / "source"
source.mkdir()
_make_hls_fixture(source)
root = tmp_path / "downloads"
root.mkdir()
config = ServiceConfig("127.0.0.1", 0, root, 1)
service = DownloadService(config)
service.update_settings({"proxy": ""})
service.start()
try:
with _http_server(source) as base_url:
task = service.create_task({"title": "collection", "directory": str(root), "payload": json.dumps([_item(base_url)])})
final = _wait_for_terminal(service, task["id"])
assert final["status"] == "completed", final
item = final["items"][0]
assert item["status"] == "completed", item
output = Path(item["output_path"])
assert output == root / "collection" / "episode.mp4"
assert output.is_file()
streams = subprocess.run(
["ffprobe", "-v", "error", "-show_streams", "-of", "json", str(output)],
capture_output=True,
text=True,
check=True,
).stdout
assert '"attached_pic": 1' in streams
finally:
service.stop()
def test_tasks_and_settings_are_not_retained_after_restart(tmp_path: Path) -> None:
root = tmp_path / "downloads"
root.mkdir()
config = ServiceConfig("127.0.0.1", 0, root, 1)
first = DownloadService(config)
first.update_settings({"default_title_template": "batch", "proxy": ""})
first.create_task({"title": "batch", "directory": str(root), "payload": json.dumps([_item("http://127.0.0.1")])})
restarted = DownloadService(config)
assert restarted.list_tasks() == []
assert restarted.bootstrap()["settings"]["default_title_template"] == DATE_TEMPLATE
assert restarted.bootstrap()["settings"]["proxy"] == DEFAULT_PROXY
assert restarted.bootstrap()["settings"]["max_active_downloads"] == 1
def test_runner_runs_multiple_videos_from_one_task_concurrently(tmp_path: Path) -> None:
store = TaskState(tmp_path, max_active_downloads=2)
runner = TaskRunner(store, worker_capacity=2, ffmpeg="ffmpeg")
lock = threading.Lock()
both_started = threading.Event()
release = threading.Event()
active_downloads = 0
peak_downloads = 0
def fake_download(task: dict[str, object], item: dict[str, object]) -> None:
nonlocal active_downloads, peak_downloads
with lock:
active_downloads += 1
peak_downloads = max(peak_downloads, active_downloads)
if active_downloads == 2:
both_started.set()
release.wait(timeout=5)
store.update_item(int(item["id"]), status="completed", stage="Completed", completed_at="now")
with lock:
active_downloads -= 1
runner._download_and_finalize = fake_download # type: ignore[method-assign]
payload = [_item("http://127.0.0.1"), {**_item("http://127.0.0.1"), "code": "episode-2", "title": "episode-2"}]
task_id = store.create_task(
title="batch",
base_dir=tmp_path,
output_dir=tmp_path / "batch",
proxy="",
items=decode_download_items(json.dumps(payload)),
)
runner.start()
try:
assert both_started.wait(timeout=3)
release.set()
deadline = time.monotonic() + 3
while time.monotonic() < deadline:
task = store.get_task(task_id)
assert task is not None
if task["status"] == "completed":
break
time.sleep(0.05)
else:
raise AssertionError("Task did not complete")
finally:
release.set()
runner.stop()
assert peak_downloads == 2
def test_cover_download_prefers_img2_s1080_url(tmp_path: Path) -> None:
store = TaskState(tmp_path)
runner = TaskRunner(store, ffmpeg="ffmpeg")
source_url = "https://cdn.example.test/img2/s720/cover.jpg?token=abc"
expected_url = "https://cdn.example.test/img2/s1080/cover.jpg?token=abc"
downloaded = tmp_path / "cover.jpg"
downloaded.write_bytes(b"cover")
requested_urls: list[str] = []
def fake_download(url: str, output_dir: Path, item_id: int, proxy: str) -> Path:
requested_urls.append(url)
return downloaded
runner._download_cover_url = fake_download # type: ignore[method-assign]
assert runner._download_cover(source_url, tmp_path, 2, "", task_id=1) == downloaded
assert requested_urls == [expected_url]
assert store.logs("warning") == []
def test_cover_download_falls_back_to_original_img2_url_and_warns(tmp_path: Path) -> None:
store = TaskState(tmp_path)
runner = TaskRunner(store, ffmpeg="ffmpeg")
source_url = "https://cdn.example.test/img2/s720/cover.jpg"
preferred_url = "https://cdn.example.test/img2/s1080/cover.jpg"
downloaded = tmp_path / "cover.jpg"
downloaded.write_bytes(b"cover")
requested_urls: list[str] = []
def fake_download(url: str, output_dir: Path, item_id: int, proxy: str) -> Path:
requested_urls.append(url)
if url == preferred_url:
raise OSError("1080p image unavailable")
return downloaded
runner._download_cover_url = fake_download # type: ignore[method-assign]
assert runner._download_cover(source_url, tmp_path, 2, "", task_id=1) == downloaded
assert requested_urls == [preferred_url, source_url]
warning = store.logs("warning")
assert len(warning) == 1
assert warning[0]["task_id"] == 1
assert warning[0]["item_id"] == 2
assert "falling back to the original URL" in warning[0]["message"]
def test_cover_download_keeps_original_url_without_img2_size_path(tmp_path: Path) -> None:
store = TaskState(tmp_path)
runner = TaskRunner(store, ffmpeg="ffmpeg")
source_url = "https://cdn.example.test/cover.jpg?path=/img2/s720/"
downloaded = tmp_path / "cover.jpg"
downloaded.write_bytes(b"cover")
requested_urls: list[str] = []
def fake_download(url: str, output_dir: Path, item_id: int, proxy: str) -> Path:
requested_urls.append(url)
return downloaded
runner._download_cover_url = fake_download # type: ignore[method-assign]
assert runner._download_cover(source_url, tmp_path, 2, "", task_id=1) == downloaded
assert requested_urls == [source_url]
def test_cli_uses_max_active_downloads() -> None:
args = build_argument_parser().parse_args(["--max-active-downloads", "3"])
assert config_from_args(args).max_active_downloads == 3
def _get(base: str, path: str) -> tuple[int, dict[str, object]]:
with urlopen(f"{base}{path}") as response:
return response.status, json.loads(response.read())
def _request(base: str, path: str, payload: dict[str, object], method: str = "POST") -> tuple[int, dict[str, object]]:
request = Request(
f"{base}{path}",
data=json.dumps(payload).encode(),
method=method,
headers={"Content-Type": "application/json"},
)
try:
with urlopen(request) as response:
return response.status, json.loads(response.read())
except Exception as error:
response = error
return response.code, json.loads(response.read()) # type: ignore[attr-defined]
def _wait_for_terminal(service: DownloadService, task_id: int, timeout: float = 30) -> dict[str, object]:
deadline = time.monotonic() + timeout
while time.monotonic() < deadline:
task = service.get_task(task_id)
assert task is not None
if task["status"] in {"completed", "partial", "failed", "cancelled"}:
return task
time.sleep(0.1)
raise AssertionError("Timed out waiting for download task")
def _make_hls_fixture(directory: Path) -> None:
subprocess.run(
[
"ffmpeg", "-hide_banner", "-loglevel", "error", "-y",
"-f", "lavfi", "-i", "testsrc2=size=160x90:rate=25:duration=1",
"-f", "lavfi", "-i", "sine=frequency=880:duration=1",
"-shortest", "-c:v", "libx264", "-pix_fmt", "yuv420p", "-c:a", "aac",
"-f", "mpegts", str(directory / "segment.ts"),
],
check=True,
)
subprocess.run(
[
"ffmpeg", "-hide_banner", "-loglevel", "error", "-y",
"-f", "lavfi", "-i", "color=c=orange:s=64x64:d=1", "-frames:v", "1", str(directory / "cover.jpg"),
],
check=True,
)
(directory / "video.m3u8").write_text("#EXTM3U\n#EXTINF:1,\nsegment.ts\n#EXT-X-ENDLIST\n")
@contextlib.contextmanager
def _http_server(directory: Path):
handler = functools.partial(http.server.SimpleHTTPRequestHandler, directory=str(directory))
server = http.server.ThreadingHTTPServer(("127.0.0.1", 0), handler)
thread = threading.Thread(target=server.serve_forever, daemon=True)
thread.start()
try:
yield f"http://127.0.0.1:{server.server_port}"
finally:
server.shutdown()
server.server_close()
thread.join()