wip: estado de trabajo pendiente antes de la vista grid (suite 230 verde)
This commit is contained in:
@@ -0,0 +1,50 @@
|
||||
"""Shared HTTP helpers for talking to YouTube CDNs.
|
||||
|
||||
Single source of truth for the headers we send to ``*.youtube.com`` endpoints
|
||||
(``ytimg.com``, ``ggpht.com``, ``youtube.com/api/timedtext``). Centralising
|
||||
them prevents the failures we hit when ``Referer: https://www.youtube.com/``
|
||||
was missing on avatar/thumbnail/subtitle fetches.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
from typing import Any, Mapping
|
||||
|
||||
import requests
|
||||
|
||||
|
||||
_YT_HEADERS: dict[str, str] = {
|
||||
# Chrome 124 on Windows. Mimics a real browser request; some Google CDN
|
||||
# endpoints (notably ``yt3.ggpht.com`` avatars and ``timedtext`` captions)
|
||||
# reject ``python-requests`` style UA-only calls with HTTP 403.
|
||||
"User-Agent": (
|
||||
"Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 "
|
||||
"(KHTML, like Gecko) Chrome/124.0 Safari/537.36"
|
||||
),
|
||||
"Referer": "https://www.youtube.com/",
|
||||
"Accept-Language": "en-US,en;q=0.9",
|
||||
# Image fetches want an Accept that lists image/* so the CDN can pick
|
||||
# the right encoded variant (avif/webp/etc).
|
||||
"Accept": (
|
||||
"image/avif,image/webp,image/png,image/jpeg,image/*,*/*;q=0.8"
|
||||
),
|
||||
}
|
||||
|
||||
|
||||
def yt_headers(extra: Mapping[str, str] | None = None) -> dict[str, str]:
|
||||
"""Return a copy of the YouTube CDN headers merged with ``extra``."""
|
||||
if extra:
|
||||
return {**_YT_HEADERS, **dict(extra)}
|
||||
return dict(_YT_HEADERS)
|
||||
|
||||
|
||||
def yt_get(url: str, *, timeout: float = 20.0,
|
||||
headers: Mapping[str, str] | None = None,
|
||||
params: Mapping[str, Any] | None = None) -> requests.Response:
|
||||
"""GET ``url`` with the shared YouTube CDN headers.
|
||||
|
||||
Centralising this lets every call site (subtitle download, video
|
||||
thumbnail cache, channel avatar cache) inherit header updates from a
|
||||
single place.
|
||||
"""
|
||||
return requests.get(url, params=params, headers=yt_headers(headers),
|
||||
timeout=timeout)
|
||||
+158
-18
@@ -3,6 +3,7 @@ from __future__ import annotations
|
||||
import logging
|
||||
import shutil
|
||||
import sys
|
||||
import time
|
||||
from pathlib import Path
|
||||
from types import SimpleNamespace
|
||||
|
||||
@@ -11,13 +12,13 @@ from rich.console import Console
|
||||
from rich.progress import Progress, SpinnerColumn, TextColumn, BarColumn, TaskProgressColumn, TimeRemainingColumn
|
||||
from rich.table import Table
|
||||
|
||||
from .config import Config, load_config
|
||||
from .config import Config, load_config, parse_languages
|
||||
from .store import Store, VideoRef, VideoRow
|
||||
from .discover import discover_channel
|
||||
from .discover import discover_channel, discover_incremental
|
||||
from .chapters import align_chapters, chapters_from_info, Chapter, Section
|
||||
from .parse import Segment
|
||||
from .render import build_filename_stem, render_markdown
|
||||
from .ratelimit import polite_sleep
|
||||
from .ratelimit import ThrottleGuard, configure_global_pacer, polite_sleep
|
||||
from .pipeline import process_video
|
||||
from .cookies import auto_import_dir, resolve_active_path
|
||||
|
||||
@@ -51,6 +52,7 @@ def cli(ctx, config_path, channel, all_channels, cookies, cookies_from_browser,
|
||||
cfg = load_config(cfg_path) if cfg_path.exists() else Config()
|
||||
if channel:
|
||||
cfg.channel_url = channel
|
||||
configure_global_pacer(cfg.delay.min_request_interval)
|
||||
store = Store(cfg.database_path_resolved)
|
||||
# import any loose cookie files into the vault (idempotent)
|
||||
try:
|
||||
@@ -72,7 +74,10 @@ def cli(ctx, config_path, channel, all_channels, cookies, cookies_from_browser,
|
||||
def _apply_filters(refs, since, no_shorts, no_live, min_duration, limit):
|
||||
filtered = list(refs)
|
||||
if since:
|
||||
filtered = [r for r in filtered if (r.upload_date or "") >= since.replace("-", "")]
|
||||
# Undated entries survive: yt-dlp's flat listing omits upload_date, so a
|
||||
# plain `>= since` would silently drop every video discovery found.
|
||||
cutoff = since.replace("-", "")
|
||||
filtered = [r for r in filtered if not r.upload_date or r.upload_date >= cutoff]
|
||||
if no_shorts:
|
||||
filtered = [r for r in filtered if "/shorts/" not in (r.url or "")]
|
||||
if no_live:
|
||||
@@ -95,20 +100,24 @@ def _apply_filters(refs, since, no_shorts, no_live, min_duration, limit):
|
||||
@click.option("--resume/--no-resume", default=True)
|
||||
@click.option("--dry-run", is_flag=True, default=False)
|
||||
@click.option("--reset-errors", is_flag=True, default=False)
|
||||
@click.option("--full", is_flag=True, default=False,
|
||||
help="Recorrer el canal entero en vez de solo los videos nuevos")
|
||||
@click.pass_obj
|
||||
def scrape(obj, limit, since, languages, no_auto, no_shorts, include_shorts, no_live, resume, dry_run, reset_errors):
|
||||
def scrape(obj, limit, since, languages, no_auto, no_shorts, include_shorts, no_live, resume, dry_run, reset_errors, full):
|
||||
"""Scrape completo: discovery + extraccion + markdown."""
|
||||
_run_scrape(obj, limit, since, languages, no_auto, no_shorts, include_shorts, no_live, resume, dry_run, reset_errors)
|
||||
_run_scrape(obj, limit, since, languages, no_auto, no_shorts, include_shorts, no_live, resume, dry_run, reset_errors, full)
|
||||
|
||||
|
||||
def _run_scrape(obj, limit, since, languages, no_auto, no_shorts, include_shorts, no_live, resume, dry_run, reset_errors):
|
||||
def _run_scrape(obj, limit, since, languages, no_auto, no_shorts, include_shorts, no_live, resume, dry_run, reset_errors, full=False):
|
||||
cfg: Config = obj.cfg
|
||||
store: Store = obj.store
|
||||
if not cfg.channel_url:
|
||||
console.print("[red]Error:[/red] falta la URL del canal. Usa --channel o config.yaml")
|
||||
sys.exit(1)
|
||||
if languages:
|
||||
cfg.languages = [l.strip() for l in languages.split(",") if l.strip()]
|
||||
# CLI flag accepts legacy CSV; per-language mode requires editing YAML.
|
||||
list_value = [l.strip() for l in languages.split(",") if l.strip()]
|
||||
cfg.languages = parse_languages(list_value, cfg.prefer_manual)
|
||||
if no_auto:
|
||||
cfg.prefer_manual = True
|
||||
if include_shorts:
|
||||
@@ -119,22 +128,52 @@ def _run_scrape(obj, limit, since, languages, no_auto, no_shorts, include_shorts
|
||||
cfg.include_live = False
|
||||
|
||||
console.print(f"[cyan]Canal:[/cyan] {cfg.channel_url}\n[cyan]DB:[/cyan] {cfg.database_path_resolved}")
|
||||
console.print("\n[bold blue]Paso 1:[/bold blue] Discovery")
|
||||
|
||||
# Incremental by default: only pull pages newer than the last video we have.
|
||||
known_channel = _known_channel_for(store, cfg.channel_url)
|
||||
incremental = cfg.sync.incremental and not full and bool(known_channel)
|
||||
known = store.known_video_ids(known_channel) if incremental else set()
|
||||
watermark = store.latest_upload_date(known_channel) if incremental else None
|
||||
|
||||
mode = (
|
||||
f"incremental (desde {_fmt_date(watermark)} en adelante)" if watermark
|
||||
else "incremental (solo videos nuevos)" if incremental
|
||||
else "completo"
|
||||
)
|
||||
console.print(f"\n[bold blue]Paso 1:[/bold blue] Discovery — [dim]{mode}[/dim]")
|
||||
try:
|
||||
with Progress(SpinnerColumn(), TextColumn("[progress.description]{task.description}"), transient=True) as prog:
|
||||
task = prog.add_task("Descubriendo videos...", total=None)
|
||||
channel_id, channel_name, refs = discover_channel(cfg.channel_url, sleep_subrequests=cfg.yt_dlp.sleep_subrequests)
|
||||
result = discover_incremental(
|
||||
cfg.channel_url, known,
|
||||
sleep_subrequests=cfg.yt_dlp.sleep_subrequests,
|
||||
window=cfg.sync.window, max_window=cfg.sync.max_window,
|
||||
overlap=cfg.sync.overlap, since=watermark,
|
||||
)
|
||||
prog.update(task, completed=1, total=1)
|
||||
except Exception as exc:
|
||||
console.print(f"[red]Error en discovery:[/red] {exc}")
|
||||
sys.exit(1)
|
||||
|
||||
store.upsert_channel(channel_id, _extract_handle(cfg.channel_url), channel_name, len(refs))
|
||||
console.print(f"[green]Canal:[/green] {channel_name} ({channel_id}) — {len(refs)} videos")
|
||||
channel_id, channel_name, refs = result.channel_id, result.channel_name, result.refs
|
||||
if result.full_scan:
|
||||
store.upsert_channel(channel_id, _extract_handle(cfg.channel_url), channel_name, len(refs))
|
||||
else:
|
||||
store.update_channel_meta(channel_id, name=channel_name)
|
||||
console.print(
|
||||
f"[green]Canal:[/green] {channel_name} ({channel_id}) — "
|
||||
f"{result.fetched} videos leidos en {result.passes} pasada(s), {result.new_count} nuevos"
|
||||
)
|
||||
if not result.caught_up:
|
||||
console.print(
|
||||
"[yellow]Aviso:[/yellow] se alcanzo el limite de ventana "
|
||||
f"({cfg.sync.max_window}). Usa --full si faltan videos antiguos."
|
||||
)
|
||||
|
||||
refs = _apply_filters(refs, since, not cfg.include_shorts, not cfg.include_live, cfg.min_duration_sec, limit)
|
||||
console.print(f"[yellow]Tras filtros:[/yellow] {len(refs)} videos")
|
||||
store.upsert_videos(refs)
|
||||
store.mark_channel_synced(channel_id)
|
||||
|
||||
if reset_errors:
|
||||
n = store.reset_errors(channel_id)
|
||||
@@ -145,8 +184,13 @@ def _run_scrape(obj, limit, since, languages, no_auto, no_shorts, include_shorts
|
||||
return
|
||||
|
||||
pending = store.get_pending(channel_id) if resume else store.get_all(channel_id)
|
||||
ref_ids = {r.video_id for r in refs}
|
||||
pending = [p for p in pending if p.video_id in ref_ids] if ref_ids else pending
|
||||
if result.full_scan:
|
||||
ref_ids = {r.video_id for r in refs}
|
||||
pending = [p for p in pending if p.video_id in ref_ids] if ref_ids else pending
|
||||
else:
|
||||
# Discovery only saw the newest slice; keep the older backlog reachable
|
||||
# but put this run's videos first so --limit still means "los mas nuevos".
|
||||
pending = _order_pending(pending, refs)
|
||||
if limit:
|
||||
pending = pending[:limit]
|
||||
if not pending:
|
||||
@@ -162,17 +206,64 @@ def _run_scrape(obj, limit, since, languages, no_auto, no_shorts, include_shorts
|
||||
with Progress(SpinnerColumn(), TextColumn("[progress.description]{task.description}"),
|
||||
BarColumn(), TaskProgressColumn(), TimeRemainingColumn()) as progress:
|
||||
task = progress.add_task("Procesando", total=len(pending))
|
||||
guard = ThrottleGuard(
|
||||
threshold=cfg.delay.throttle_threshold,
|
||||
base=cfg.delay.backoff_base,
|
||||
cap=cfg.delay.backoff_cap,
|
||||
)
|
||||
for i, row in enumerate(pending):
|
||||
progress.update(task, description=f"{row.video_id} {(row.title or '')[:30]}", completed=i)
|
||||
process_video(row, cfg, store, channel_name, channel_id, cfg.channel_url,
|
||||
cookies_file=cookie_path, cookies_from_browser=obj.cookies_from_browser)
|
||||
status = process_video(row, cfg, store, channel_name, channel_id, cfg.channel_url,
|
||||
cookies_file=cookie_path, cookies_from_browser=obj.cookies_from_browser)
|
||||
progress.advance(task)
|
||||
if status == "done":
|
||||
guard.note_success()
|
||||
else:
|
||||
failed = store.get_video(row.video_id)
|
||||
wait = guard.note_failure(getattr(failed, "error_msg", None) or status)
|
||||
if guard.tripped:
|
||||
# Everything still pending stays pending: that is what makes
|
||||
# this recoverable instead of 500 rows marked failed.
|
||||
console.print(f"\n[bold red]Detenido:[/bold red] {guard.tripped_reason}")
|
||||
console.print(f"[yellow]{len(pending) - i - 1} videos sin tocar, siguen pendientes.[/yellow]")
|
||||
break
|
||||
if wait > 0:
|
||||
console.print(f"[yellow]YouTube nos esta limitando; esperando {wait:.1f}s[/yellow]")
|
||||
time.sleep(wait)
|
||||
if i < len(pending) - 1:
|
||||
polite_sleep(cfg.delay.min_seconds, cfg.delay.max_seconds)
|
||||
console.print()
|
||||
_print_stats(store, channel_id)
|
||||
|
||||
|
||||
def _known_channel_for(store: Store, channel_url: str) -> str | None:
|
||||
"""Match a channel URL against an already-tracked channel, by handle or id.
|
||||
|
||||
Without a hit there is no local history to stop at, so discovery has to walk
|
||||
the whole channel — which is correct for a first run.
|
||||
"""
|
||||
if not channel_url:
|
||||
return None
|
||||
handle = _extract_handle(channel_url).lstrip("@").lower()
|
||||
tail = channel_url.rstrip("/").split("/")[-1]
|
||||
for ch in store.list_channels():
|
||||
cid = ch.get("channel_id") or ""
|
||||
ch_handle = (ch.get("handle") or "").lstrip("@").lower()
|
||||
if handle and ch_handle == handle:
|
||||
return cid
|
||||
if cid and (cid == tail or f"/channel/{cid}" in channel_url):
|
||||
return cid
|
||||
return None
|
||||
|
||||
|
||||
def _order_pending(rows: list, refs: list) -> list:
|
||||
"""Videos from this run's window first, then the rest of the backlog."""
|
||||
by_id = {r.video_id: r for r in rows}
|
||||
ordered = [by_id.pop(r.video_id) for r in refs if r.video_id in by_id]
|
||||
ordered.extend(by_id.values())
|
||||
return ordered
|
||||
|
||||
|
||||
def _print_dry_run(refs):
|
||||
table = Table(show_lines=False)
|
||||
table.add_column("Fecha", style="dim")
|
||||
@@ -186,6 +277,55 @@ def _print_dry_run(refs):
|
||||
console.print(f"... y {len(refs) - 50} mas")
|
||||
|
||||
|
||||
# --------------------------------------------------------------------------- recovery
|
||||
|
||||
@cli.command("reset")
|
||||
@click.option("--status", "statuses", multiple=True,
|
||||
type=click.Choice(Store.RETRYABLE_STATUSES),
|
||||
help="Estados a reintentar (por defecto: error y no_subtitles)")
|
||||
@click.option("--yes", is_flag=True, default=False, help="No preguntar")
|
||||
@click.pass_obj
|
||||
def reset_cmd(obj, statuses, yes):
|
||||
"""Devolver videos atascados en error/no_subtitles a 'pending' para reintentarlos."""
|
||||
store: Store = obj.store
|
||||
targets = tuple(statuses) if statuses else Store.RETRYABLE_STATUSES
|
||||
channels = _channel_targets(obj)
|
||||
channel_id = channels[0] if len(channels) == 1 else None
|
||||
|
||||
counts = store.retryable_counts(channel_id)
|
||||
affected = sum(counts.get(s, 0) for s in targets)
|
||||
if not affected:
|
||||
console.print("[green]No hay videos que reintentar.[/green]")
|
||||
return
|
||||
detail = ", ".join(f"{s}={counts.get(s, 0)}" for s in targets)
|
||||
scope = channel_id or "todos los canales"
|
||||
if not yes and not click.confirm(f"Reintentar {affected} videos ({detail}) en {scope}?"):
|
||||
return
|
||||
n = store.reset_videos(channel_id, targets)
|
||||
console.print(f"[yellow]{n}[/yellow] videos vueltos a 'pending'.")
|
||||
|
||||
|
||||
@cli.command("reconcile")
|
||||
@click.option("--prune", is_flag=True, default=False,
|
||||
help="Ademas, devolver a 'pending' los 'done' cuyo .md ya no existe")
|
||||
@click.pass_obj
|
||||
def reconcile_cmd(obj, prune):
|
||||
"""Re-escanear data/markdown y hacer que la DB coincida con el disco."""
|
||||
from .segments import reconcile_markdown
|
||||
cfg: Config = obj.cfg
|
||||
store: Store = obj.store
|
||||
result = reconcile_markdown(
|
||||
store, Path(cfg.output_dir_resolved),
|
||||
log=lambda m: console.print(f"[dim]{m}[/dim]"), prune=prune,
|
||||
)
|
||||
table = Table(title="Reconciliacion disco <-> DB")
|
||||
table.add_column("Metrica", style="bold")
|
||||
table.add_column("Valor", justify="right")
|
||||
for k, v in result.items():
|
||||
table.add_row(k, str(v))
|
||||
console.print(table)
|
||||
|
||||
|
||||
# --------------------------------------------------------------------------- search
|
||||
|
||||
@cli.command("search")
|
||||
@@ -313,8 +453,8 @@ def channels_add(obj, url):
|
||||
cfg: Config = obj.cfg
|
||||
cfg.channel_url = url
|
||||
console.print("[cyan]Resolviendo canal...[/cyan]")
|
||||
channel_id, name, refs = discover_channel(url, sleep_subrequests=cfg.yt_dlp.sleep_subrequests)
|
||||
obj.store.upsert_channel(channel_id, _extract_handle(url), name, len(refs))
|
||||
channel_id, name, avatar, refs = discover_channel(url, sleep_subrequests=cfg.yt_dlp.sleep_subrequests)
|
||||
obj.store.upsert_channel(channel_id, _extract_handle(url), name, len(refs), avatar=avatar)
|
||||
obj.store.upsert_videos(refs)
|
||||
console.print(f"[green]Added:[/green] {name} ({channel_id}) — {len(refs)} videos")
|
||||
|
||||
|
||||
@@ -9,28 +9,79 @@ import yaml
|
||||
|
||||
@dataclass
|
||||
class DelayConfig:
|
||||
"""Pacing and failure handling for everything that touches YouTube.
|
||||
|
||||
`min_seconds`/`max_seconds` are the randomised gap between units of work
|
||||
(one video, one channel). `backoff_*` and `throttle_threshold` only come
|
||||
into play once YouTube starts refusing: the backoff pair feeds the
|
||||
truncated-exponential-with-jitter formula Google documents for its own
|
||||
APIs, and the threshold is how many consecutive rate-limit responses end
|
||||
the run instead of burning through the rest of the queue.
|
||||
"""
|
||||
|
||||
min_seconds: float = 1.5
|
||||
max_seconds: float = 3.5
|
||||
backoff_base: float = 2.0
|
||||
backoff_cap: float = 60.0
|
||||
throttle_threshold: int = 3
|
||||
# Minimum gap between *any* two requests to YouTube from this process,
|
||||
# including the ones the `/api/tools/*` endpoints make outside the job
|
||||
# runner. 0 disables the shared pacer and leaves only the per-unit sleeps.
|
||||
min_request_interval: float = 0.0
|
||||
# Bytes/sec ceiling for audio downloads (yt-dlp `ratelimit`). 0 = unlimited.
|
||||
audio_rate_limit: int = 0
|
||||
|
||||
|
||||
@dataclass
|
||||
class YtDlpConfig:
|
||||
retries: int = 10
|
||||
# Seconds between the individual HTTP requests inside one extraction.
|
||||
# Passed to yt-dlp as `sleep_interval_requests`; the name here is the
|
||||
# project's own and predates the discovery that the option this used to be
|
||||
# forwarded under (`sleep_subrequests`) does not exist in yt-dlp at all.
|
||||
sleep_subrequests: float = 2.0
|
||||
extractor_retries: int = 3
|
||||
socket_timeout: float = 30.0
|
||||
|
||||
|
||||
@dataclass
|
||||
class SyncConfig:
|
||||
"""Incremental channel sync — how far back a routine re-scan looks.
|
||||
|
||||
The /videos tab is reverse-chronological and every extra page is another
|
||||
request to YouTube, so a sync fetches `window` entries and stops as soon as
|
||||
it has seen `overlap` consecutive videos already in the DB. Only if the
|
||||
whole window turns out to be new does it widen (doubling up to `max_window`),
|
||||
which is the case where the channel really did publish a lot since last time.
|
||||
"""
|
||||
|
||||
incremental: bool = True
|
||||
window: int = 30
|
||||
max_window: int = 300
|
||||
overlap: int = 3
|
||||
|
||||
|
||||
@dataclass
|
||||
class Config:
|
||||
channel_url: str = ""
|
||||
languages: list[str] = field(default_factory=lambda: ["es", "en"])
|
||||
# Per-language subtitle preference. Each value is one of:
|
||||
# "manual" - only manually uploaded captions
|
||||
# "auto" - only YouTube-auto-generated captions
|
||||
# "any" - defer to the legacy `prefer_manual` flag
|
||||
# Legacy form: ``["es", "en"]`` - treated as ``{lang: <prefer_manual_default>}``.
|
||||
# Default to "any" (manual first, auto as fallback). A manual-only default
|
||||
# silently yields nothing on the many channels that publish only
|
||||
# auto-generated captions, and stores that as `no_subtitles`.
|
||||
languages: dict[str, str] = field(
|
||||
default_factory=lambda: {"es": "any", "en": "any"}
|
||||
)
|
||||
prefer_manual: bool = True
|
||||
include_shorts: bool = False
|
||||
include_live: bool = True
|
||||
min_duration_sec: int = 0
|
||||
delay: DelayConfig = field(default_factory=DelayConfig)
|
||||
yt_dlp: YtDlpConfig = field(default_factory=YtDlpConfig)
|
||||
sync: SyncConfig = field(default_factory=SyncConfig)
|
||||
database_path: str = "data/state.db"
|
||||
output_dir: str = "data/markdown"
|
||||
template_path: str = "templates/video.md.j2"
|
||||
@@ -49,6 +100,27 @@ class Config:
|
||||
return Path(self.template_path).resolve()
|
||||
|
||||
|
||||
# Modes accepted in the ``languages`` dict. Kept here so tests and CLI
|
||||
# share a single definition without importing the private extract constant.
|
||||
LANGUAGE_MODES = ("manual", "auto", "any")
|
||||
|
||||
|
||||
def parse_languages(raw: Any, prefer_manual: bool) -> dict[str, str]:
|
||||
"""Normalise legacy list / new dict / None into ``{lang: mode}``."""
|
||||
default = "manual" if prefer_manual else "auto"
|
||||
if raw is None:
|
||||
return {}
|
||||
if isinstance(raw, dict):
|
||||
out: dict[str, str] = {}
|
||||
for lang, mode in raw.items():
|
||||
m = str(mode).lower().strip()
|
||||
out[str(lang)] = m if m in LANGUAGE_MODES else "any"
|
||||
return out
|
||||
if isinstance(raw, (list, tuple)):
|
||||
return {str(l): default for l in raw}
|
||||
return {}
|
||||
|
||||
|
||||
def load_config(path: str | Path) -> Config:
|
||||
p = Path(path)
|
||||
if not p.exists():
|
||||
@@ -61,10 +133,15 @@ def load_config(path: str | Path) -> Config:
|
||||
def _build_config(raw: dict[str, Any]) -> Config:
|
||||
delay_raw = raw.get("delay") or {}
|
||||
ydl_raw = raw.get("yt_dlp") or {}
|
||||
sync_raw = raw.get("sync") or {}
|
||||
prefer_manual = bool(raw.get("prefer_manual", True))
|
||||
languages = parse_languages(raw.get("languages"), prefer_manual)
|
||||
if not languages:
|
||||
languages = {"es": "any", "en": "any"}
|
||||
return Config(
|
||||
channel_url=raw.get("channel_url", ""),
|
||||
languages=list(raw.get("languages", ["es", "en"])),
|
||||
prefer_manual=bool(raw.get("prefer_manual", True)),
|
||||
languages=languages,
|
||||
prefer_manual=prefer_manual,
|
||||
include_shorts=bool(raw.get("include_shorts", False)),
|
||||
include_live=bool(raw.get("include_live", True)),
|
||||
min_duration_sec=int(raw.get("min_duration_sec", 0)),
|
||||
@@ -73,10 +150,21 @@ def _build_config(raw: dict[str, Any]) -> Config:
|
||||
max_seconds=float(delay_raw.get("max_seconds", 3.5)),
|
||||
backoff_base=float(delay_raw.get("backoff_base", 2.0)),
|
||||
backoff_cap=float(delay_raw.get("backoff_cap", 60.0)),
|
||||
throttle_threshold=max(1, int(delay_raw.get("throttle_threshold", 3))),
|
||||
min_request_interval=max(0.0, float(delay_raw.get("min_request_interval", 0.0))),
|
||||
audio_rate_limit=max(0, int(delay_raw.get("audio_rate_limit", 0))),
|
||||
),
|
||||
yt_dlp=YtDlpConfig(
|
||||
retries=int(ydl_raw.get("retries", 10)),
|
||||
sleep_subrequests=float(ydl_raw.get("sleep_subrequests", 2.0)),
|
||||
extractor_retries=max(0, int(ydl_raw.get("extractor_retries", 3))),
|
||||
socket_timeout=float(ydl_raw.get("socket_timeout", 30.0)),
|
||||
),
|
||||
sync=SyncConfig(
|
||||
incremental=bool(sync_raw.get("incremental", True)),
|
||||
window=max(1, int(sync_raw.get("window", 30))),
|
||||
max_window=max(1, int(sync_raw.get("max_window", 300))),
|
||||
overlap=max(1, int(sync_raw.get("overlap", 3))),
|
||||
),
|
||||
database_path=raw.get("database_path", "data/state.db"),
|
||||
output_dir=raw.get("output_dir", "data/markdown"),
|
||||
|
||||
+230
-34
@@ -2,12 +2,13 @@ from __future__ import annotations
|
||||
|
||||
import logging
|
||||
from dataclasses import dataclass
|
||||
from typing import Any
|
||||
from typing import Any, Mapping
|
||||
|
||||
import requests
|
||||
import yt_dlp
|
||||
|
||||
from ._yt_http import yt_get
|
||||
from .parse import Segment, parse_auto_dump
|
||||
from .ratelimit import GLOBAL_PACER, ydl_throttle_opts
|
||||
|
||||
log = logging.getLogger(__name__)
|
||||
|
||||
@@ -26,75 +27,199 @@ class VideoData:
|
||||
segments: list[Segment]
|
||||
subtitle: SubtitlePick | None
|
||||
has_chapters: bool
|
||||
# Why `segments` came back empty. "No transcript" has several very
|
||||
# different causes — the video genuinely has no captions, the language
|
||||
# policy rejected the tracks that do exist, or the download was throttled —
|
||||
# and collapsing them into one terminal status hides recoverable failures.
|
||||
skip_reason: str | None = None
|
||||
|
||||
|
||||
_LANGUAGE_MODES = ("manual", "auto", "any")
|
||||
|
||||
|
||||
def _sources_for(mode: str, prefer_manual: bool, manual: dict, auto: dict) -> list[tuple[str, dict]]:
|
||||
"""Return ordered list of (label, tracks_dict) to try for ``mode``.
|
||||
|
||||
``any`` defers to the legacy :data:`prefer_manual` global default.
|
||||
``manual`` / ``auto`` force the track family even if the other has
|
||||
a higher-priority language elsewhere in the iteration.
|
||||
"""
|
||||
if mode == "manual":
|
||||
return [("manual", manual)]
|
||||
if mode == "auto":
|
||||
return [("auto", auto)]
|
||||
if prefer_manual:
|
||||
return [("manual", manual), ("auto", auto)]
|
||||
return [("auto", auto), ("manual", manual)]
|
||||
|
||||
|
||||
def extract_video(
|
||||
video_url: str,
|
||||
languages: list[str],
|
||||
languages: Mapping[str, str] | list[str],
|
||||
retries: int = 10,
|
||||
sleep_subrequests: float = 2.0,
|
||||
prefer_manual: bool = True,
|
||||
cookies_file: str | None = None,
|
||||
cookies_from_browser: str | None = None,
|
||||
extractor_retries: int = 3,
|
||||
socket_timeout: float = 30.0,
|
||||
) -> VideoData:
|
||||
"""Run yt-dlp on ``video_url`` and pull the preferred subtitle track.
|
||||
|
||||
``languages`` may be either a list (legacy, every entry uses
|
||||
``prefer_manual``) or a mapping ``{lang: mode}`` where ``mode`` is one
|
||||
of ``"manual"``, ``"auto"`` or ``"any"``. The mapping form is the
|
||||
preferred interface because it lets you mix per-language policies such
|
||||
as ``{"en": "manual", "es": "auto", "pt": "any"}``.
|
||||
"""
|
||||
languages_dict = _coerce_languages(languages, prefer_manual)
|
||||
|
||||
ydl_opts: dict[str, Any] = {
|
||||
"writesubtitles": True,
|
||||
"writeautomaticsub": True,
|
||||
"subtitleslangs": languages,
|
||||
"subtitleslangs": list(languages_dict.keys()),
|
||||
"skip_download": True,
|
||||
"quiet": True,
|
||||
"no_warnings": True,
|
||||
"retries": retries,
|
||||
"sleep_subrequests": sleep_subrequests,
|
||||
"noprogress": True,
|
||||
# Never let yt-dlp probe formats: it costs one HTTP request per format
|
||||
# and we only ever want captions and metadata.
|
||||
"check_formats": None,
|
||||
**ydl_throttle_opts(
|
||||
sleep_subrequests,
|
||||
extractor_retries=extractor_retries,
|
||||
socket_timeout=socket_timeout,
|
||||
),
|
||||
}
|
||||
if cookies_file:
|
||||
ydl_opts["cookiefile"] = cookies_file
|
||||
if cookies_from_browser:
|
||||
ydl_opts["cookiesfrombrowser"] = (cookies_from_browser,)
|
||||
|
||||
# One video extraction is two requests: the watch page and the InnerTube
|
||||
# player call. yt-dlp spaces them itself; the pacer needs to know they exist.
|
||||
GLOBAL_PACER.wait(cost=2)
|
||||
with yt_dlp.YoutubeDL(ydl_opts) as ydl:
|
||||
info = ydl.extract_info(video_url, download=False)
|
||||
|
||||
pick = pick_subtitle(info, languages, prefer_manual)
|
||||
pick = pick_subtitle(info, languages_dict, prefer_manual)
|
||||
segments: list[Segment] = []
|
||||
skip_reason: str | None = None
|
||||
if pick:
|
||||
raw = _download_subtitle(pick.url)
|
||||
raw, dl_error = _download_subtitle(pick.url)
|
||||
if raw:
|
||||
segments = parse_auto_dump(raw)
|
||||
if not segments:
|
||||
log.warning("Could not parse subtitle for %s (format=%s)", video_url, pick.ext)
|
||||
skip_reason = f"subtitle downloaded but parsed empty (lang={pick.lang}, format={pick.ext})"
|
||||
else:
|
||||
skip_reason = (
|
||||
f"subtitle track found (lang={pick.lang}, {pick.source}) but the download failed "
|
||||
f"— usually throttling; retry later [{dl_error}]"
|
||||
)
|
||||
else:
|
||||
skip_reason = describe_missing_subtitle(info, languages_dict)
|
||||
|
||||
has_chapters = bool(info.get("chapters"))
|
||||
return VideoData(info=info, segments=segments, subtitle=pick, has_chapters=has_chapters)
|
||||
return VideoData(
|
||||
info=info, segments=segments, subtitle=pick,
|
||||
has_chapters=has_chapters, skip_reason=skip_reason,
|
||||
)
|
||||
|
||||
|
||||
def pick_subtitle(info: dict[str, Any], languages: list[str], prefer_manual: bool = True) -> SubtitlePick | None:
|
||||
def describe_missing_subtitle(info: dict[str, Any], languages: Mapping[str, str]) -> str:
|
||||
"""Explain why no track matched, distinguishing 'none exist' from 'policy rejected them'.
|
||||
|
||||
A channel that only publishes auto-generated captions scanned under a
|
||||
manual-only policy yields nothing — which is a config problem, not a
|
||||
property of the video, and the message has to say so.
|
||||
"""
|
||||
manual = {k: v for k, v in (info.get("subtitles") or {}).items() if v}
|
||||
auto = {k: v for k, v in (info.get("automatic_captions") or {}).items() if v}
|
||||
if not manual and not auto:
|
||||
return "no caption tracks published for this video"
|
||||
|
||||
wanted = ", ".join(f"{lang}={mode}" for lang, mode in languages.items()) or "(none configured)"
|
||||
modes = {str(m).lower() for m in languages.values()}
|
||||
parts = [f"no track matched the language policy ({wanted})"]
|
||||
parts.append(f"available: {len(manual)} manual, {len(auto)} auto")
|
||||
if auto and not manual and modes == {"manual"}:
|
||||
parts.append(
|
||||
"this video has ONLY auto-generated captions — set the language mode "
|
||||
"to 'any' or 'auto' to use them"
|
||||
)
|
||||
return "; ".join(parts)
|
||||
|
||||
|
||||
def _coerce_languages(languages: Mapping[str, str] | list[str] | None,
|
||||
prefer_manual: bool) -> dict[str, str]:
|
||||
"""Normalise legacy list / new dict / None into ``{lang: mode}``."""
|
||||
default = "manual" if prefer_manual else "auto"
|
||||
if languages is None:
|
||||
return {}
|
||||
if isinstance(languages, Mapping):
|
||||
out: dict[str, str] = {}
|
||||
for lang, mode in languages.items():
|
||||
m = str(mode).lower().strip()
|
||||
if m not in _LANGUAGE_MODES:
|
||||
m = "any"
|
||||
out[str(lang)] = m
|
||||
return out
|
||||
if isinstance(languages, (list, tuple)):
|
||||
return {str(l): default for l in languages}
|
||||
return {}
|
||||
|
||||
|
||||
def pick_subtitle(info: dict[str, Any],
|
||||
languages: Mapping[str, str],
|
||||
prefer_manual: bool = True) -> SubtitlePick | None:
|
||||
"""Pick the best subtitle track for ``info`` honouring per-language mode.
|
||||
|
||||
See :func:`extract_video` for the ``languages`` schema. ``prefer_manual``
|
||||
is only consulted for entries whose mode is ``"any"``.
|
||||
"""
|
||||
manual = info.get("subtitles") or {}
|
||||
auto = info.get("automatic_captions") or {}
|
||||
|
||||
ordered_sources: list[tuple[str, dict[str, Any]]]
|
||||
if prefer_manual:
|
||||
ordered_sources = [("manual", manual), ("auto", auto)]
|
||||
else:
|
||||
ordered_sources = [("auto", auto), ("manual", manual)]
|
||||
if isinstance(languages, (list, tuple)):
|
||||
# legacy path: convert on the fly
|
||||
default = "manual" if prefer_manual else "auto"
|
||||
languages = {l: default for l in languages}
|
||||
|
||||
for source_label, tracks in ordered_sources:
|
||||
for lang in languages:
|
||||
# Config order is a preference between languages we can read, not an
|
||||
# instruction to accept a machine translation when the real transcript is
|
||||
# sitting right there. An English channel scanned under {es, es-419, en}
|
||||
# was yielding Spanish auto-translations of English speech.
|
||||
#
|
||||
# Two passes rather than a reorder. The reorder alone needed to know the
|
||||
# spoken language, and when neither an `-orig` key nor `info["language"]`
|
||||
# was present it silently fell back to config order and reintroduced the
|
||||
# bug. Rejecting translations outright in the first pass needs no such
|
||||
# knowledge: whatever language it lands on, it is the one actually spoken.
|
||||
ordered = list(languages.items())
|
||||
spoken = original_language(info)
|
||||
if spoken and any(_normalize_lang(l) == spoken for l, _ in ordered):
|
||||
ordered.sort(key=lambda kv: _normalize_lang(kv[0]) != spoken)
|
||||
|
||||
for allow_translations in (False, True):
|
||||
for lang, mode in ordered:
|
||||
if mode not in _LANGUAGE_MODES:
|
||||
mode = "any"
|
||||
sources = _sources_for(mode, prefer_manual, manual, auto)
|
||||
normalized = _normalize_lang(lang)
|
||||
for track_lang, formats in tracks.items():
|
||||
if _normalize_lang(track_lang) != normalized:
|
||||
continue
|
||||
if not formats:
|
||||
continue
|
||||
pick = _pick_best_format(formats)
|
||||
if pick:
|
||||
return SubtitlePick(
|
||||
url=pick["url"],
|
||||
ext=pick["ext"],
|
||||
lang=track_lang,
|
||||
source="manual" if source_label == "manual" else "auto",
|
||||
)
|
||||
for source_label, tracks in sources:
|
||||
for track_lang, formats in _ordered_tracks(tracks, normalized):
|
||||
if not allow_translations and _is_translation(track_lang, formats):
|
||||
continue
|
||||
pick = _pick_best_format(formats)
|
||||
if pick:
|
||||
return SubtitlePick(
|
||||
url=pick["url"],
|
||||
ext=pick["ext"],
|
||||
lang=track_lang,
|
||||
source=source_label,
|
||||
)
|
||||
return None
|
||||
|
||||
|
||||
@@ -115,11 +240,82 @@ def _normalize_lang(code: str) -> str:
|
||||
return base
|
||||
|
||||
|
||||
def _download_subtitle(url: str) -> str | None:
|
||||
def _is_original_track(code: str) -> bool:
|
||||
"""True for YouTube's original-ASR track, which it suffixes with ``-orig``.
|
||||
|
||||
YouTube publishes the speech-recognised track as ``<lang>-orig`` and then a
|
||||
long tail of machine translations keyed by bare language code — including a
|
||||
translation *into the video's own language*. So on a Spanish video both
|
||||
``es-orig`` and ``es`` exist, and only the first is the real transcript.
|
||||
"""
|
||||
return code.replace("_", "-").lower().endswith("-orig")
|
||||
|
||||
|
||||
def original_language(info: dict[str, Any]) -> str | None:
|
||||
"""The language actually spoken in the video, normalised, or None.
|
||||
|
||||
Prefers the ``-orig`` track that YouTube itself publishes over ``info`` keys,
|
||||
because the ``-orig`` suffix is direct evidence from the caption list while
|
||||
``language`` is metadata that YouTube localises along with the title.
|
||||
"""
|
||||
for code in (info.get("automatic_captions") or {}):
|
||||
if _is_original_track(code):
|
||||
return _normalize_lang(code)
|
||||
for key in ("language", "original_language"):
|
||||
val = info.get(key)
|
||||
if isinstance(val, str) and val:
|
||||
return _normalize_lang(val)
|
||||
return None
|
||||
|
||||
|
||||
def _is_translation(code: str, formats: list[dict[str, Any]] | None) -> bool:
|
||||
"""True when this track is YouTube machine-translating some other track.
|
||||
|
||||
The caption URL says so outright: yt-dlp builds a translated track by
|
||||
appending ``tlang=`` to the base track's URL and omits it when the target
|
||||
equals the source language. That is direct evidence, unlike the ``-orig``
|
||||
naming convention, and it is what lets us reject a translation even for a
|
||||
video whose spoken language we could not otherwise determine.
|
||||
"""
|
||||
if _is_original_track(code):
|
||||
return False
|
||||
for fmt in formats or []:
|
||||
url = fmt.get("url") or ""
|
||||
if "tlang=" in url:
|
||||
return True
|
||||
return False
|
||||
|
||||
|
||||
def _ordered_tracks(tracks: dict, normalized: str) -> list[tuple[str, Any]]:
|
||||
"""Tracks matching `normalized`, original-ASR first.
|
||||
|
||||
Without this the picker took whichever key yt-dlp happened to list first.
|
||||
That silently returned the right thing on Spanish channels (``es-orig``
|
||||
sorts before ``es``) and the wrong thing everywhere else.
|
||||
"""
|
||||
matches = [(c, f) for c, f in tracks.items() if _normalize_lang(c) == normalized and f]
|
||||
matches.sort(key=lambda kv: not _is_original_track(kv[0]))
|
||||
return matches
|
||||
|
||||
|
||||
def _download_subtitle(url: str, *, timeout: float = 15.0) -> tuple[str | None, str | None]:
|
||||
"""Fetch a caption track. Returns (text, error_description).
|
||||
|
||||
The error text is returned rather than only logged because a 429 here is
|
||||
how YouTube throttling most often shows up on this path, and the circuit
|
||||
breaker upstream can only see it if it survives into the stored reason.
|
||||
"""
|
||||
try:
|
||||
resp = requests.get(url, timeout=15, headers={"User-Agent": "Mozilla/5.0"})
|
||||
GLOBAL_PACER.wait()
|
||||
resp = yt_get(url, timeout=timeout)
|
||||
resp.raise_for_status()
|
||||
return resp.text
|
||||
except requests.RequestException as exc:
|
||||
return resp.text, None
|
||||
except Exception as exc: # pylint: disable=broad-except
|
||||
log.error("Failed to download subtitle from %s: %s", url, exc)
|
||||
return None
|
||||
# Lead with a normalised "HTTP Error <status>" token. Downstream, the
|
||||
# throttle detector has to recognise a 429 here, and depending on the
|
||||
# prose is fragile: a 429 served without a reason phrase (routine over
|
||||
# HTTP/2) says nothing about "too many requests".
|
||||
status = getattr(getattr(exc, "response", None), "status_code", None)
|
||||
prefix = f"HTTP Error {status}: " if status else ""
|
||||
return None, f"{prefix}{type(exc).__name__}: {exc}"
|
||||
|
||||
@@ -6,9 +6,9 @@ from typing import Callable
|
||||
|
||||
from .config import Config
|
||||
from .cookies import resolve_active_path
|
||||
from .discover import discover_channel
|
||||
from .discover import discover_incremental
|
||||
from .pipeline import process_video
|
||||
from .ratelimit import polite_sleep
|
||||
from .ratelimit import ThrottleGuard, polite_sleep
|
||||
from .store import Store
|
||||
|
||||
|
||||
@@ -65,13 +65,34 @@ def _run_once(
|
||||
if on_log:
|
||||
on_log(msg)
|
||||
|
||||
channel_id_found, channel_name, refs = discover_channel(
|
||||
cfg.channel_url, sleep_subrequests=cfg.yt_dlp.sleep_subrequests
|
||||
# A watch loop re-runs on an interval; walking the whole channel every tick
|
||||
# is exactly what burns through YouTube's tolerance. Only look at what is
|
||||
# newer than the last video we already have.
|
||||
known = store.known_video_ids(channel_id) if (channel_id and cfg.sync.incremental) else set()
|
||||
result = discover_incremental(
|
||||
cfg.channel_url,
|
||||
known,
|
||||
sleep_subrequests=cfg.yt_dlp.sleep_subrequests,
|
||||
window=cfg.sync.window,
|
||||
max_window=cfg.sync.max_window,
|
||||
overlap=cfg.sync.overlap,
|
||||
since=store.latest_upload_date(channel_id) if known else None,
|
||||
keep=_keep_ref(cfg),
|
||||
)
|
||||
target_channel = channel_id or channel_id_found
|
||||
store.upsert_channel(target_channel, _extract_handle(cfg.channel_url), channel_name, len(refs))
|
||||
channel_name, refs = result.channel_name, result.refs
|
||||
target_channel = channel_id or result.channel_id
|
||||
if result.full_scan:
|
||||
store.upsert_channel(
|
||||
target_channel, _extract_handle(cfg.channel_url), channel_name, len(refs), avatar=result.avatar
|
||||
)
|
||||
else:
|
||||
store.update_channel_meta(target_channel, name=channel_name, avatar=result.avatar)
|
||||
store.upsert_videos(refs)
|
||||
_emit(f"watch: discovered {len(refs)} videos on {channel_name}")
|
||||
store.mark_channel_synced(target_channel)
|
||||
_emit(
|
||||
f"watch: scanned {result.fetched} newest video(s) on {channel_name} — "
|
||||
f"{result.new_count} new"
|
||||
)
|
||||
|
||||
pending = store.get_pending(target_channel)
|
||||
if not pending:
|
||||
@@ -86,16 +107,48 @@ def _run_once(
|
||||
cookie_path = vault
|
||||
_emit(f"watch: using active vault cookie {vault}")
|
||||
|
||||
guard = ThrottleGuard(
|
||||
threshold=cfg.delay.throttle_threshold,
|
||||
base=cfg.delay.backoff_base,
|
||||
cap=cfg.delay.backoff_cap,
|
||||
)
|
||||
for row in pending:
|
||||
process_video(
|
||||
status = process_video(
|
||||
row, cfg, store, channel_name, target_channel, cfg.channel_url,
|
||||
cookies_file=cookie_path,
|
||||
cookies_from_browser=cookies_from_browser,
|
||||
on_log=on_log,
|
||||
)
|
||||
if status == "done":
|
||||
guard.note_success()
|
||||
else:
|
||||
failed = store.get_video(row.video_id)
|
||||
wait = guard.note_failure(getattr(failed, "error_msg", None) or status)
|
||||
if guard.tripped:
|
||||
# Unattended loop: stop this pass rather than spend the whole
|
||||
# backlog against a throttled session. The next tick retries.
|
||||
_emit(f"watch: {guard.tripped_reason} — stopping this pass")
|
||||
return
|
||||
if wait > 0:
|
||||
_emit(f"watch: throttled, waiting {wait:.1f}s")
|
||||
time.sleep(wait)
|
||||
polite_sleep(cfg.delay.min_seconds, cfg.delay.max_seconds)
|
||||
|
||||
|
||||
def _keep_ref(cfg: Config):
|
||||
"""Same shorts/live rules the store was populated under — see jobs._keep_ref."""
|
||||
|
||||
def keep(r) -> bool:
|
||||
url = r.url or ""
|
||||
if not cfg.include_shorts and "/shorts/" in url:
|
||||
return False
|
||||
if not cfg.include_live and url.startswith("https://www.youtube.com/live/"):
|
||||
return False
|
||||
return True
|
||||
|
||||
return keep
|
||||
|
||||
|
||||
def _extract_handle(url: str) -> str:
|
||||
if "@" in url:
|
||||
return "@" + url.split("@", 1)[1].split("/", 1)[0]
|
||||
|
||||
+130
-37
@@ -8,13 +8,84 @@ from typing import Callable
|
||||
from .chapters import align_chapters, chapters_from_info
|
||||
from .config import Config
|
||||
from .extract import extract_video
|
||||
from .ratelimit import is_rate_limited
|
||||
from .render import build_filename_stem, render_markdown
|
||||
from .store import Store, VideoRow
|
||||
from ._yt_http import yt_get
|
||||
|
||||
|
||||
log = logging.getLogger(__name__)
|
||||
|
||||
|
||||
def _render_and_retire(
|
||||
cfg: Config,
|
||||
video_id: str,
|
||||
channel_name: str,
|
||||
context: dict,
|
||||
previous: str | None,
|
||||
) -> tuple[Path, str]:
|
||||
"""Write the note for `video_id` and delete whatever file it used to own.
|
||||
|
||||
Identity is the video id; the *filename* is derived from the title, and
|
||||
titles are not stable. YouTube serves them localised — the same video came
|
||||
back as "La controversia de Claude Fable 5" on one pass and "The Claude
|
||||
Fable controversy 5" on the next — and creators rename videos outright.
|
||||
Re-rendering under a new stem without retiring the old path leaves an
|
||||
orphan the database no longer references, which is how a library ends up
|
||||
with more notes than videos.
|
||||
|
||||
Shared by `process_video` and `re_render_videos` because they previously
|
||||
each built the stem their own way and drifted: one used the raw compact
|
||||
date, the other the normalised one, and the result was 94 files for 61 rows.
|
||||
"""
|
||||
stem = build_filename_stem(
|
||||
upload_date=context["upload_date"],
|
||||
title=context["title"],
|
||||
template=cfg.filename_template,
|
||||
video_id=video_id,
|
||||
)
|
||||
out_subdir = Path(cfg.output_dir_resolved) / _safe_dirname(channel_name)
|
||||
md_path = render_markdown(cfg.template_path_resolved, out_subdir, stem, context)
|
||||
|
||||
out_root = Path(cfg.output_dir_resolved).parent
|
||||
try:
|
||||
rel = md_path.relative_to(out_root) if md_path.is_relative_to(out_root) else md_path
|
||||
except ValueError:
|
||||
rel = md_path
|
||||
rel_str = str(rel)
|
||||
|
||||
if previous and previous != rel_str:
|
||||
old = out_root / previous
|
||||
try:
|
||||
if old.exists() and old.resolve() != md_path.resolve():
|
||||
old.unlink()
|
||||
log.info("retired superseded note %s -> %s", previous, rel_str)
|
||||
except OSError as exc:
|
||||
log.warning("could not remove superseded note %s: %s", old, exc)
|
||||
return md_path, rel_str
|
||||
|
||||
|
||||
def _store_metadata(store: Store, row: VideoRow, info: dict) -> None:
|
||||
"""Persist everything the extraction learned that is not the transcript.
|
||||
|
||||
Split out of the render path so it can run before the no-transcript exit:
|
||||
these fields are already paid for by the time we know whether captions
|
||||
downloaded, and `update_video_metadata` only writes the keys it is given.
|
||||
"""
|
||||
store.update_video_metadata(
|
||||
row.video_id,
|
||||
view_count=info.get("view_count"),
|
||||
like_count=info.get("like_count"),
|
||||
tags=info.get("tags") or None,
|
||||
thumbnail=info.get("thumbnail"),
|
||||
description=(info.get("description") or "").strip() or None,
|
||||
)
|
||||
# discovery in flat mode reports no upload_date, so this is where it lands
|
||||
ud = info.get("upload_date") or row.upload_date
|
||||
if ud:
|
||||
store.set_upload_date(row.video_id, ud.replace("-", "") if "-" in str(ud) else str(ud))
|
||||
|
||||
|
||||
def process_video(
|
||||
row: VideoRow,
|
||||
cfg: Config,
|
||||
@@ -42,22 +113,41 @@ def process_video(
|
||||
prefer_manual=cfg.prefer_manual,
|
||||
cookies_file=cookies_file,
|
||||
cookies_from_browser=cookies_from_browser,
|
||||
extractor_retries=cfg.yt_dlp.extractor_retries,
|
||||
socket_timeout=cfg.yt_dlp.socket_timeout,
|
||||
)
|
||||
except Exception as exc:
|
||||
log.error("Error extrayendo %s: %s", row.video_id, exc)
|
||||
store.mark_error(row.video_id, str(exc))
|
||||
return "error"
|
||||
|
||||
info = data.info
|
||||
store.set_availability(row.video_id, info.get("availability"))
|
||||
|
||||
# Metadata is persisted BEFORE the no-transcript exit. The extraction already
|
||||
# cost its requests and the info dict is in hand; discarding it because the
|
||||
# separate caption fetch failed means a retry re-spends them for data we
|
||||
# already had. Measured after a throttling incident: five rows left with
|
||||
# upload_date, view_count, description and thumbnail all NULL.
|
||||
_store_metadata(store, row, info)
|
||||
|
||||
if not data.segments:
|
||||
log.warning("Sin transcripcion para %s", row.video_id)
|
||||
store.mark_status(row.video_id, "no_subtitles")
|
||||
reason = data.skip_reason or "no transcript"
|
||||
log.warning("Sin transcripcion para %s: %s", row.video_id, reason)
|
||||
_emit(f"{row.video_id}: {reason}")
|
||||
# A throttled caption fetch is not "this video has no captions". Both
|
||||
# states are retryable, but only `error` is honest about the cause, and
|
||||
# `no_subtitles` counts are what tell you a channel publishes none.
|
||||
# Which of the three requests YouTube refused should not decide this.
|
||||
if is_rate_limited(reason):
|
||||
store.mark_error(row.video_id, reason)
|
||||
return "error"
|
||||
store.mark_status(row.video_id, "no_subtitles", reason)
|
||||
return "no_subtitles"
|
||||
|
||||
chapters = chapters_from_info(data.info)
|
||||
sections = align_chapters(data.segments, chapters)
|
||||
|
||||
info = data.info
|
||||
|
||||
# persist segments + rich metadata to DB (for search, stats, webapp)
|
||||
store.store_segments(row.video_id, data.segments)
|
||||
seg_json = json.dumps(
|
||||
@@ -70,20 +160,10 @@ def process_video(
|
||||
)
|
||||
store.update_video_metadata(
|
||||
row.video_id,
|
||||
view_count=info.get("view_count"),
|
||||
like_count=info.get("like_count"),
|
||||
tags=info.get("tags") or None,
|
||||
thumbnail=info.get("thumbnail"),
|
||||
description=(info.get("description") or "").strip() or None,
|
||||
chapters_json=ch_json,
|
||||
segments_json=seg_json,
|
||||
)
|
||||
|
||||
# ensure upload_date is populated (discovery sometimes lacks it)
|
||||
ud = info.get("upload_date") or row.upload_date
|
||||
if ud:
|
||||
store.set_upload_date(row.video_id, ud.replace("-", "") if "-" in str(ud) else str(ud))
|
||||
|
||||
context = {
|
||||
"video_id": row.video_id,
|
||||
"title": info.get("title") or row.title or row.video_id,
|
||||
@@ -103,22 +183,12 @@ def process_video(
|
||||
"sections": sections,
|
||||
}
|
||||
|
||||
stem = build_filename_stem(
|
||||
upload_date=context["upload_date"],
|
||||
title=context["title"],
|
||||
template=cfg.filename_template,
|
||||
md_path, rel_str = _render_and_retire(
|
||||
cfg, row.video_id, channel_name, context, row.markdown_path
|
||||
)
|
||||
out_subdir = Path(cfg.output_dir_resolved) / _safe_dirname(channel_name)
|
||||
md_path = render_markdown(cfg.template_path_resolved, out_subdir, stem, context)
|
||||
|
||||
out_root = Path(cfg.output_dir_resolved).parent
|
||||
try:
|
||||
rel = md_path.relative_to(out_root) if md_path.is_relative_to(out_root) else md_path
|
||||
except ValueError:
|
||||
rel = md_path
|
||||
store.mark_done(
|
||||
row.video_id,
|
||||
str(rel),
|
||||
rel_str,
|
||||
data.subtitle.lang if data.subtitle else None,
|
||||
data.subtitle.source if data.subtitle else None,
|
||||
data.has_chapters,
|
||||
@@ -150,20 +220,38 @@ def thumbnail_url_for(video_row) -> str:
|
||||
|
||||
def cache_thumbnail(store, video_id: str, thumbnails_dir) -> bool:
|
||||
"""Download + cache a video's thumbnail to <thumbnails_dir>/<video_id>.jpg. Returns True on success/existing."""
|
||||
import requests as _requests
|
||||
out = Path(thumbnails_dir) / f"{video_id}.jpg"
|
||||
if out.exists():
|
||||
if out.exists() and out.stat().st_size > 0:
|
||||
return True
|
||||
v = store.get_video(video_id)
|
||||
url = thumbnail_url_for(v) if v else f"https://i.ytimg.com/vi/{video_id}/hqdefault.jpg"
|
||||
try:
|
||||
r = _requests.get(url, timeout=15, headers={"User-Agent": "Mozilla/5.0"})
|
||||
r = yt_get(url, timeout=20.0)
|
||||
r.raise_for_status()
|
||||
out.parent.mkdir(parents=True, exist_ok=True)
|
||||
out.write_bytes(r.content)
|
||||
return True
|
||||
if r.content:
|
||||
out.parent.mkdir(parents=True, exist_ok=True)
|
||||
out.write_bytes(r.content)
|
||||
return True
|
||||
except Exception:
|
||||
return False
|
||||
return False
|
||||
|
||||
|
||||
def cache_channel_avatar(channel_id: str, url: str, avatars_dir) -> bool:
|
||||
"""Download + cache a channel avatar to <avatars_dir>/<channel_id>.jpg. Returns True on success/existing."""
|
||||
out = Path(avatars_dir) / f"{channel_id}.jpg"
|
||||
out.parent.mkdir(parents=True, exist_ok=True)
|
||||
if out.exists() and out.stat().st_size > 0:
|
||||
return True
|
||||
try:
|
||||
r = yt_get(url, timeout=20.0)
|
||||
r.raise_for_status()
|
||||
if r.content:
|
||||
out.write_bytes(r.content)
|
||||
return True
|
||||
except Exception:
|
||||
pass
|
||||
return False
|
||||
|
||||
|
||||
def re_render_videos(store: Store, cfg: Config, channel_id: str | None = None) -> int:
|
||||
@@ -193,14 +281,19 @@ def re_render_videos(store: Store, cfg: Config, channel_id: str | None = None) -
|
||||
context = {
|
||||
"video_id": v.video_id, "title": v.title or v.video_id,
|
||||
"channel_name": ch.get("name") or "", "channel_id": v.channel_id, "channel_url": "",
|
||||
"upload_date": v.upload_date or "", "duration": v.duration or 0, "url": v.url,
|
||||
# Must match process_video's normalisation (pipeline.py:96): the
|
||||
# frontmatter is a bidirectional contract that backfill re-reads.
|
||||
"upload_date": _normalize_date(v.upload_date), "duration": v.duration or 0, "url": v.url,
|
||||
"transcript_lang": v.transcript_lang or "", "transcript_src": v.transcript_src or "",
|
||||
"view_count": v.view_count, "like_count": v.like_count, "tags": tags,
|
||||
"thumbnail": v.thumbnail or "", "description": v.description or "", "sections": sections,
|
||||
}
|
||||
stem = build_filename_stem(v.upload_date, v.title or v.video_id, cfg.filename_template)
|
||||
out_subdir = Path(cfg.output_dir_resolved) / _safe_dirname(ch.get("name") or "unknown")
|
||||
render_markdown(cfg.template_path_resolved, out_subdir, stem, context)
|
||||
md_path, rel_str = _render_and_retire(
|
||||
cfg, v.video_id, ch.get("name") or "unknown", context, v.markdown_path
|
||||
)
|
||||
# mark_done is the only writer of markdown_path; without it the DB
|
||||
# keeps pointing at the pre-render file.
|
||||
store.mark_done(v.video_id, rel_str, v.transcript_lang, v.transcript_src, bool(chapters))
|
||||
n += 1
|
||||
except Exception as exc:
|
||||
log.warning("re-render failed for %s: %s", v.video_id, exc)
|
||||
|
||||
+278
-4
@@ -1,14 +1,288 @@
|
||||
"""Politeness primitives shared by every path that reaches YouTube.
|
||||
|
||||
Three separate concerns live here, and they are not interchangeable:
|
||||
|
||||
* `Pacer` spaces requests out so a burst never leaves this process. It is
|
||||
process-global on purpose — the JobManager serialises *jobs*, but the
|
||||
`/api/tools/*` endpoints run outside it, so without a shared pacer two
|
||||
consumers can hammer YouTube while each believes it is being polite.
|
||||
* `backoff_delay` is what to wait *after* a failure. It follows the algorithm
|
||||
Google documents for its own APIs: `min(base * 2**n + jitter, cap)`, with the
|
||||
jitter redrawn on every attempt so retries from concurrent clients do not
|
||||
re-synchronise into waves.
|
||||
* `ThrottleGuard` decides when to stop trying. YouTube's throttle is a session
|
||||
ban of up to an hour; once it lands, every further request is both useless
|
||||
and harmful. The guard is what turns "511 videos marked failed" into "stopped
|
||||
after 3, kept them retryable".
|
||||
|
||||
Deliberately absent: any retry of the request itself. yt-dlp's YouTube
|
||||
extractor explicitly excludes 403/429 from its own RetryManager, so a throttled
|
||||
call fails once and comes back here — retrying it in a tight loop is the exact
|
||||
behaviour that earns the ban in the first place.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import random
|
||||
import re
|
||||
import threading
|
||||
import time
|
||||
from dataclasses import dataclass, field
|
||||
from typing import Callable
|
||||
|
||||
# Substrings that mean "YouTube is refusing because of *rate*, not content".
|
||||
#
|
||||
# Kept distinct from `Store.PERMANENT_ERROR_PATTERNS`: those describe videos we
|
||||
# will never get (members-only, deleted). These describe videos we could get if
|
||||
# we asked more slowly, so they must stay retryable and must trip the breaker.
|
||||
_THROTTLE_SIGNS = (
|
||||
"rate-limited by youtube",
|
||||
"this content isn't available, try again later",
|
||||
"sign in to confirm you're not a bot",
|
||||
"sign in to confirm your age", # same bot-wall, different copy
|
||||
"http error 429",
|
||||
"too many requests",
|
||||
"ratelimitexceeded",
|
||||
"userratelimitexceeded",
|
||||
"temporarily blocked",
|
||||
)
|
||||
|
||||
# 403 with a quota reason is NOT the same failure: Google documents it as
|
||||
# daily and explicitly warns retries may not work for hours (AIP-194). Treated
|
||||
# as fatal-for-now rather than as something to back off and retry into.
|
||||
_QUOTA_SIGNS = (
|
||||
"quotaexceeded",
|
||||
"dailylimitexceeded",
|
||||
)
|
||||
|
||||
# Two spellings reach us and they are not the same string:
|
||||
# yt-dlp -> "HTTP Error 429: Too Many Requests"
|
||||
# requests -> "HTTPError: 429 Client Error: Too Many Requests for url: ..."
|
||||
# The original pattern only matched the first, so the second was detected purely
|
||||
# by its "too many requests" prose — and a 429 served without a reason phrase
|
||||
# (routine over HTTP/2) carries no such prose and slipped through the breaker.
|
||||
_HTTP_STATUS = re.compile(
|
||||
r"HTTP\s*Error[:\s]+(\d{3})"
|
||||
r"|(\d{3})\s+(?:Client|Server)\s+Error",
|
||||
re.IGNORECASE,
|
||||
)
|
||||
|
||||
#: Statuses Google documents as retryable, and which mean "slow down" here.
|
||||
_RETRYABLE_STATUSES = frozenset({"408", "429"})
|
||||
|
||||
|
||||
def is_rate_limited(error: object) -> bool:
|
||||
"""True when `error` looks like YouTube throttling us rather than a bad video."""
|
||||
text = str(error or "").lower()
|
||||
if any(sign in text for sign in _THROTTLE_SIGNS):
|
||||
return True
|
||||
for m in _HTTP_STATUS.finditer(text):
|
||||
status = m.group(1) or m.group(2)
|
||||
if status in _RETRYABLE_STATUSES:
|
||||
return True
|
||||
return False
|
||||
|
||||
|
||||
def is_quota_exhausted(error: object) -> bool:
|
||||
"""True for Google's daily-quota refusals, which backing off will not fix."""
|
||||
text = str(error or "").lower()
|
||||
return any(sign in text for sign in _QUOTA_SIGNS)
|
||||
|
||||
|
||||
def polite_sleep(min_s: float = 1.5, max_s: float = 3.5) -> None:
|
||||
delay = random.uniform(min_s, max_s)
|
||||
time.sleep(delay)
|
||||
"""Randomised pause between units of work (one video, one channel)."""
|
||||
if max_s < min_s:
|
||||
max_s = min_s
|
||||
time.sleep(random.uniform(min_s, max_s))
|
||||
|
||||
|
||||
def backoff_sleep(attempt: int, base: float = 2.0, cap: float = 60.0) -> None:
|
||||
delay = min(cap, base * (2 ** attempt) + random.uniform(0, 1))
|
||||
def backoff_delay(attempt: int, base: float = 2.0, cap: float = 60.0) -> float:
|
||||
"""Truncated exponential backoff with full jitter, per Google's retry guidance.
|
||||
|
||||
`attempt` is 0-indexed, so the first wait after a failure is ~`base`.
|
||||
Returns the delay instead of sleeping so callers can log it and tests can
|
||||
assert on it without spending real seconds.
|
||||
"""
|
||||
if attempt < 0:
|
||||
attempt = 0
|
||||
# 2**attempt overflows into absurd floats long before it matters; clamp the
|
||||
# exponent so a runaway counter can't turn into an OverflowError.
|
||||
exponent = min(attempt, 32)
|
||||
return min(cap, base * (2**exponent) + random.uniform(0, 1))
|
||||
|
||||
|
||||
def backoff_sleep(attempt: int, base: float = 2.0, cap: float = 60.0) -> float:
|
||||
"""Sleep `backoff_delay(...)` and return how long it waited."""
|
||||
delay = backoff_delay(attempt, base, cap)
|
||||
time.sleep(delay)
|
||||
return delay
|
||||
|
||||
|
||||
class Pacer:
|
||||
"""Enforces a minimum gap between requests, process-wide and thread-safe.
|
||||
|
||||
`wait()` blocks only for the remainder of the gap, so a slow caller never
|
||||
pays twice: if the previous request already took longer than `min_interval`
|
||||
there is nothing left to wait for.
|
||||
"""
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
min_interval: float = 0.0,
|
||||
*,
|
||||
clock: Callable[[], float] = time.monotonic,
|
||||
sleeper: Callable[[float], None] = time.sleep,
|
||||
) -> None:
|
||||
self.min_interval = max(0.0, float(min_interval))
|
||||
self._clock = clock
|
||||
self._sleeper = sleeper
|
||||
self._lock = threading.Lock()
|
||||
self._next_at = 0.0
|
||||
|
||||
def configure(self, min_interval: float) -> None:
|
||||
with self._lock:
|
||||
self.min_interval = max(0.0, float(min_interval))
|
||||
|
||||
def wait(self, cost: int = 1) -> float:
|
||||
"""Block until the next request is allowed. Returns seconds actually slept.
|
||||
|
||||
`cost` is how many HTTP requests the caller is about to make. It matters
|
||||
because the caller is usually yt-dlp: one `extract_info` on a video is
|
||||
two requests (watch page + InnerTube player), and charging it as one
|
||||
made the pacer under-count by half. yt-dlp spaces those two internally
|
||||
via `sleep_interval_requests`; `cost` is what keeps this pacer's idea of
|
||||
the budget honest about them.
|
||||
"""
|
||||
cost = max(1, int(cost))
|
||||
with self._lock:
|
||||
if self.min_interval <= 0:
|
||||
return 0.0
|
||||
now = self._clock()
|
||||
delay = self._next_at - now
|
||||
if delay < 0:
|
||||
delay = 0.0
|
||||
# Reserve the slots before releasing the lock so concurrent callers
|
||||
# queue up behind each other instead of all reading the same `now`.
|
||||
self._next_at = now + delay + self.min_interval * cost
|
||||
if delay > 0:
|
||||
self._sleeper(delay)
|
||||
return delay
|
||||
|
||||
def penalise(self, seconds: float) -> None:
|
||||
"""Push the next allowed request out by `seconds` (used after a 429)."""
|
||||
if seconds <= 0:
|
||||
return
|
||||
with self._lock:
|
||||
self._next_at = max(self._next_at, self._clock() + seconds)
|
||||
|
||||
|
||||
#: Shared by `_yt_http.yt_get` and every yt-dlp entry point. Idle (0.0) until
|
||||
#: something calls `configure_global_pacer`, so importing this module never
|
||||
#: slows down a test suite that does not opt in.
|
||||
GLOBAL_PACER = Pacer(0.0)
|
||||
|
||||
|
||||
def configure_global_pacer(min_interval: float) -> None:
|
||||
GLOBAL_PACER.configure(min_interval)
|
||||
|
||||
|
||||
def ydl_throttle_opts(
|
||||
sleep_requests: float = 2.0,
|
||||
*,
|
||||
extractor_retries: int = 3,
|
||||
socket_timeout: float = 30.0,
|
||||
) -> dict[str, object]:
|
||||
"""The politeness half of every `ydl_opts` dict in this project.
|
||||
|
||||
Exists because the option names are easy to get subtly wrong, and yt-dlp
|
||||
silently ignores keys it does not recognise — this project shipped
|
||||
`sleep_subrequests` (not a real option) for its entire history, so nothing
|
||||
ever slept between the sub-requests of an extraction.
|
||||
|
||||
Only `sleep_interval_requests` throttles *extraction*; `sleep_interval` and
|
||||
`max_sleep_interval` fire in the file downloader and never trigger under
|
||||
`skip_download`, which is every metadata path here.
|
||||
"""
|
||||
return {
|
||||
# Sleeps inside InfoExtractor._request_webpage, i.e. before each HTTP
|
||||
# call an extractor makes: watch page, InnerTube player/browse, and the
|
||||
# continuation pages of a channel tab.
|
||||
"sleep_interval_requests": max(0.0, float(sleep_requests)),
|
||||
# `retries` governs the downloader; extraction retries are this one.
|
||||
# Note yt-dlp's YouTube extractor refuses to retry 403/429 at all, so
|
||||
# this only covers transient 5xx and network errors.
|
||||
"extractor_retries": int(extractor_retries),
|
||||
"socket_timeout": float(socket_timeout),
|
||||
}
|
||||
|
||||
|
||||
@dataclass
|
||||
class ThrottleGuard:
|
||||
"""Circuit breaker for a batch of YouTube work.
|
||||
|
||||
Feed it every outcome. It answers one question — "should this run keep
|
||||
going?" — and, while it still says yes, how long to wait first.
|
||||
|
||||
The counter is *consecutive*: an isolated throttled video between successes
|
||||
is noise, three in a row means the session is banned and everything after
|
||||
it will fail too. Only the consecutive form distinguishes those.
|
||||
"""
|
||||
|
||||
threshold: int = 3
|
||||
base: float = 2.0
|
||||
cap: float = 60.0
|
||||
|
||||
consecutive: int = 0
|
||||
throttled_total: int = 0
|
||||
tripped_reason: str | None = None
|
||||
_waits: list[float] = field(default_factory=list)
|
||||
|
||||
@property
|
||||
def tripped(self) -> bool:
|
||||
return self.tripped_reason is not None
|
||||
|
||||
def note_success(self) -> None:
|
||||
"""A request went through: the session is healthy again."""
|
||||
self.consecutive = 0
|
||||
|
||||
def note_failure(self, error: object) -> float:
|
||||
"""Record a failed unit of work. Returns seconds the caller should wait.
|
||||
|
||||
Non-throttle failures (a private video, a parse error) reset the
|
||||
consecutive counter: they say nothing about our request rate, and
|
||||
letting them accumulate would trip the breaker on a channel that simply
|
||||
has a few dead videos.
|
||||
"""
|
||||
if is_quota_exhausted(error):
|
||||
self.throttled_total += 1
|
||||
self.tripped_reason = (
|
||||
"YouTube reports the quota is exhausted; Google documents this as "
|
||||
"daily, so retrying now cannot succeed"
|
||||
)
|
||||
return 0.0
|
||||
|
||||
if not is_rate_limited(error):
|
||||
self.consecutive = 0
|
||||
return 0.0
|
||||
|
||||
self.throttled_total += 1
|
||||
self.consecutive += 1
|
||||
if self.consecutive >= self.threshold:
|
||||
self.tripped_reason = (
|
||||
f"{self.consecutive} consecutive rate-limit responses from YouTube; "
|
||||
"the session is throttled (YouTube states up to an hour) and further "
|
||||
"requests would only mark healthy videos as failed"
|
||||
)
|
||||
return 0.0
|
||||
|
||||
delay = backoff_delay(self.consecutive - 1, self.base, self.cap)
|
||||
self._waits.append(delay)
|
||||
GLOBAL_PACER.penalise(delay)
|
||||
return delay
|
||||
|
||||
def summary(self) -> str:
|
||||
if self.tripped_reason:
|
||||
return f"stopped: {self.tripped_reason}"
|
||||
if self.throttled_total:
|
||||
return f"{self.throttled_total} throttled response(s) absorbed by backoff"
|
||||
return "no throttling seen"
|
||||
|
||||
@@ -63,10 +63,26 @@ def render_markdown(
|
||||
return out_file
|
||||
|
||||
|
||||
def build_filename_stem(upload_date: str | None, title: str, template: str = "{upload_date}_{slug}") -> str:
|
||||
def build_filename_stem(
|
||||
upload_date: str | None,
|
||||
title: str,
|
||||
template: str = "{upload_date}_{slug}",
|
||||
video_id: str | None = None,
|
||||
) -> str:
|
||||
"""Filename for a video's note, from `template`.
|
||||
|
||||
`{video_id}` is offered because the other two variables are unstable:
|
||||
YouTube serves titles localised (the same video came back Spanish on one
|
||||
pass and English on the next) and creators rename things. A template
|
||||
including the id makes the file identifiable from disk alone, without
|
||||
consulting the database. It is opt-in — the default is unchanged so
|
||||
existing libraries keep their filenames.
|
||||
"""
|
||||
from slugify import slugify
|
||||
date_part = upload_date or "unknown-date"
|
||||
slug = slugify(title, max_length=60) or "untitled"
|
||||
stem = template.format(upload_date=date_part, slug=slug, title=title)
|
||||
stem = template.format(
|
||||
upload_date=date_part, slug=slug, title=title, video_id=video_id or "",
|
||||
)
|
||||
safe = "".join(c for c in stem if c not in r'\/:*?"<>|')
|
||||
return safe.strip().strip(".")
|
||||
|
||||
+113
-1
@@ -122,7 +122,9 @@ def backfill_from_markdown(
|
||||
for md_path in md_files:
|
||||
try:
|
||||
text = md_path.read_text(encoding="utf-8")
|
||||
except OSError as exc:
|
||||
except (OSError, UnicodeDecodeError) as exc:
|
||||
# UnicodeDecodeError is a ValueError, not an OSError — letting it
|
||||
# escape aborted the loop and silently skipped every later file.
|
||||
_log(f"backfill: skip unreadable {md_path}: {exc}")
|
||||
continue
|
||||
parsed = parse_markdown(text)
|
||||
@@ -170,6 +172,116 @@ def backfill_from_markdown(
|
||||
return n
|
||||
|
||||
|
||||
def reconcile_markdown(
|
||||
store: Store,
|
||||
md_root: Path,
|
||||
log: Callable[[str], None] | None = None,
|
||||
*,
|
||||
prune: bool = False,
|
||||
) -> dict[str, int]:
|
||||
"""Make the DB agree with what is actually on disk.
|
||||
|
||||
`backfill_from_markdown` fills in segments and metadata but never touches
|
||||
`status` or `markdown_path`, so a video whose .md exists can sit at
|
||||
`error`/`no_subtitles`/`pending` forever and the UI keeps showing a failure
|
||||
for work that is already done. This walks the markdown tree and repairs:
|
||||
|
||||
- a row with a real .md but a non-done status -> marked done
|
||||
|
||||
`prune=True` additionally sends `done` rows whose .md has disappeared back
|
||||
to pending. That direction is opt-in because it is destructive when aimed
|
||||
at the wrong root: pointed at an empty or unrelated markdown tree it would
|
||||
demote every finished video in the database. It is also skipped outright
|
||||
when the tree contains no .md at all, which is never a real "everything was
|
||||
deleted" state — it means the root is wrong.
|
||||
|
||||
Returns counts so the caller can report what changed. Idempotent.
|
||||
"""
|
||||
def _log(msg: str) -> None:
|
||||
if log:
|
||||
log(msg)
|
||||
else:
|
||||
logging.getLogger(__name__).info(msg)
|
||||
|
||||
md_root = Path(md_root)
|
||||
data_root = md_root.parent
|
||||
out = {
|
||||
"scanned": 0, "repaired_done": 0, "orphan_md": 0,
|
||||
"missing_md": 0, "backfilled": 0, "stale_dupe": 0,
|
||||
}
|
||||
if not md_root.exists():
|
||||
_log(f"reconcile: markdown root not found: {md_root}")
|
||||
return out
|
||||
|
||||
out["backfilled"] = backfill_from_markdown(store, md_root, log=log)
|
||||
|
||||
seen: dict[str, Path] = {}
|
||||
for md_path in sorted(md_root.rglob("*.md")):
|
||||
out["scanned"] += 1
|
||||
try:
|
||||
text = md_path.read_text(encoding="utf-8")
|
||||
except (OSError, UnicodeDecodeError) as exc:
|
||||
_log(f"reconcile: skip unreadable {md_path}: {exc}")
|
||||
continue
|
||||
meta = parse_markdown(text).metadata
|
||||
video_id = meta.get("video_id")
|
||||
if not video_id:
|
||||
continue
|
||||
row = store.get_video(video_id)
|
||||
if not row:
|
||||
out["orphan_md"] += 1
|
||||
continue
|
||||
rel = md_path.relative_to(data_root).as_posix()
|
||||
|
||||
# A second .md for a video the DB already resolves elsewhere. Older
|
||||
# re-renders built the filename from a differently-formatted date, so
|
||||
# they wrote a sibling file the DB never learned about; it is dead
|
||||
# weight that every later scan has to wade through.
|
||||
canonical = (row.markdown_path or "").replace("\\", "/")
|
||||
if video_id in seen or (row.status == "done" and canonical and canonical != rel):
|
||||
out["stale_dupe"] += 1
|
||||
if prune:
|
||||
try:
|
||||
md_path.unlink()
|
||||
_log(f"reconcile: removed stale duplicate {rel}")
|
||||
except OSError as exc:
|
||||
_log(f"reconcile: could not remove {rel}: {exc}")
|
||||
else:
|
||||
_log(f"reconcile: stale duplicate (use prune to delete): {rel}")
|
||||
continue
|
||||
|
||||
seen[video_id] = md_path
|
||||
if row.status != "done" or not row.markdown_path:
|
||||
store.mark_done(
|
||||
video_id, rel,
|
||||
meta.get("transcript_lang") or row.transcript_lang,
|
||||
meta.get("transcript_src") or row.transcript_src,
|
||||
bool(meta.get("has_chapters")) or bool(row.has_chapters),
|
||||
)
|
||||
out["repaired_done"] += 1
|
||||
_log(f"reconcile: {video_id} had a .md on disk but status={row.status} -> done")
|
||||
|
||||
# The other direction is destructive, so it needs both an explicit opt-in
|
||||
# and evidence that we are looking at a real markdown tree.
|
||||
if prune and out["scanned"]:
|
||||
for row in store.get_all():
|
||||
if row.status != "done" or row.video_id in seen:
|
||||
continue
|
||||
path = data_root / row.markdown_path if row.markdown_path else None
|
||||
if path is None or not path.exists():
|
||||
store.mark_status(row.video_id, "pending", "markdown file missing on disk")
|
||||
out["missing_md"] += 1
|
||||
_log(f"reconcile: {row.video_id} marked done but .md is gone -> pending")
|
||||
elif prune:
|
||||
_log("reconcile: markdown tree is empty — refusing to prune (wrong root?)")
|
||||
|
||||
_log(
|
||||
"reconcile: scanned {scanned} .md, repaired {repaired_done}, "
|
||||
"re-queued {missing_md}, orphans {orphan_md}, stale duplicates {stale_dupe}".format(**out)
|
||||
)
|
||||
return out
|
||||
|
||||
|
||||
def _to_int(value: str | None) -> int | None:
|
||||
if value is None:
|
||||
return None
|
||||
|
||||
@@ -383,6 +383,24 @@ def build_router(store: Store, cfg: Config, jobs) -> APIRouter:
|
||||
raise HTTPException(404, "audio not downloaded yet")
|
||||
return FileResponse(str(p), filename=f"{video_id}.mp3", media_type="audio/mpeg")
|
||||
|
||||
@r.api_route("/videos/{video_id}/media", methods=["GET", "HEAD"])
|
||||
def video_media(video_id: str, request: Request, inline: bool = False):
|
||||
v = store.get_video(video_id)
|
||||
if not v or v.video_download_status != "done" or not v.video_path:
|
||||
raise HTTPException(404, "video not downloaded yet")
|
||||
root = Path(cfg.output_dir_resolved).parent.resolve()
|
||||
p = (root / v.video_path).resolve()
|
||||
try:
|
||||
p.relative_to(root)
|
||||
except ValueError:
|
||||
raise HTTPException(500, "invalid stored video path")
|
||||
if not p.exists():
|
||||
raise HTTPException(404, "video file missing on disk")
|
||||
return FileResponse(
|
||||
str(p), filename=v.video_filename or p.name, media_type="video/webm",
|
||||
content_disposition_type="inline" if inline else "attachment",
|
||||
)
|
||||
|
||||
# -------------------------------------------------- thumbnails (local cache)
|
||||
@r.post("/tools/thumbnails")
|
||||
def download_thumbnails(payload: dict):
|
||||
@@ -614,6 +632,22 @@ def build_router(store: Store, cfg: Config, jobs) -> APIRouter:
|
||||
job_id = jobs.enqueue(channel_id, opts)
|
||||
return {"job_id": job_id}
|
||||
|
||||
@r.post("/tools/video")
|
||||
async def tools_video(payload: dict):
|
||||
video_ids = (payload or {}).get("video_ids") or []
|
||||
if len(video_ids) != 1:
|
||||
raise HTTPException(400, "exactly one video_id required")
|
||||
video_id = str(video_ids[0])
|
||||
v = store.get_video(video_id)
|
||||
if not v:
|
||||
raise HTTPException(404, "video not found")
|
||||
md_path = Path(cfg.output_dir_resolved).parent / v.markdown_path if v.markdown_path else None
|
||||
if v.status != "done" or not v.markdown_path or not md_path or not md_path.exists():
|
||||
raise HTTPException(409, "markdown must be generated and present before downloading video")
|
||||
opts = {"mode": "video", "video_ids": [video_id]}
|
||||
job_id = jobs.enqueue(v.channel_id, opts)
|
||||
return {"job_id": job_id}
|
||||
|
||||
return r
|
||||
|
||||
|
||||
@@ -633,6 +667,10 @@ def _video_dict(v) -> dict:
|
||||
"transcript_lang": v.transcript_lang, "has_chapters": bool(v.has_chapters),
|
||||
"view_count": v.view_count, "like_count": v.like_count, "tags": tags,
|
||||
"thumbnail": v.thumbnail, "markdown_path": v.markdown_path,
|
||||
"video_download_status": v.video_download_status or "not_downloaded",
|
||||
"video_path": v.video_path, "video_filename": v.video_filename,
|
||||
"video_size": v.video_size, "video_downloaded_at": v.video_downloaded_at,
|
||||
"video_error": v.video_error,
|
||||
# The date the row was ordered by. For a video discovery has not
|
||||
# extracted yet there is no real upload_date, so this is inferred from
|
||||
# its position in the channel listing — flagged, never passed off as
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
from pathlib import Path
|
||||
|
||||
from fastapi import FastAPI
|
||||
@@ -8,7 +9,8 @@ from fastapi.staticfiles import StaticFiles
|
||||
|
||||
from ..config import Config, load_config
|
||||
from ..cookies import auto_import_dir
|
||||
from ..segments import backfill_from_markdown
|
||||
from ..ratelimit import configure_global_pacer
|
||||
from ..segments import reconcile_markdown
|
||||
from ..store import Store
|
||||
from .api import build_router
|
||||
from .jobs import JobManager
|
||||
@@ -23,18 +25,34 @@ def create_app(cfg: Config | None = None, db_path: str | Path | None = None) ->
|
||||
cfg.database_path = str(db_path)
|
||||
store = Store(cfg.database_path_resolved)
|
||||
|
||||
# The job runner serialises jobs, but /api/tools/* endpoints run outside it
|
||||
# on the threadpool. The pacer is the only thing that stops those two from
|
||||
# hitting YouTube simultaneously, so it has to be armed before any router.
|
||||
configure_global_pacer(cfg.delay.min_request_interval)
|
||||
|
||||
# auto-import any loose cookies into the vault
|
||||
try:
|
||||
auto_import_dir(store)
|
||||
except Exception:
|
||||
pass
|
||||
# auto-backfill segments from existing markdown (idempotent)
|
||||
# uvicorn only configures its own loggers, so the root logger has no
|
||||
# handler and everything yt_scraper logs at startup vanishes.
|
||||
pkg_log = logging.getLogger("yt_scraper")
|
||||
if not pkg_log.handlers:
|
||||
handler = logging.StreamHandler()
|
||||
handler.setFormatter(logging.Formatter("%(levelname)s %(name)s: %(message)s"))
|
||||
pkg_log.addHandler(handler)
|
||||
pkg_log.setLevel(logging.INFO)
|
||||
|
||||
# Reconcile, not just backfill: backfill_from_markdown populates segments
|
||||
# and metadata but never touches `status`/`markdown_path`, so a video whose
|
||||
# .md is already on disk would keep showing a failure after every restart.
|
||||
try:
|
||||
md_root = cfg.output_dir_resolved
|
||||
if Path(md_root).exists():
|
||||
backfill_from_markdown(store, Path(md_root))
|
||||
md_root = Path(cfg.output_dir_resolved)
|
||||
if md_root.exists():
|
||||
reconcile_markdown(store, md_root, log=pkg_log.info)
|
||||
except Exception:
|
||||
pass
|
||||
pkg_log.exception("startup reconcile failed")
|
||||
|
||||
app = FastAPI(title="yt-scraper platform", version="1.0.0")
|
||||
|
||||
|
||||
+488
-36
@@ -6,14 +6,15 @@ import threading
|
||||
import time
|
||||
import uuid
|
||||
from collections import deque
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
|
||||
from ..config import Config, load_config
|
||||
from ..config import Config, load_config, parse_languages
|
||||
from ..cookies import resolve_active_path
|
||||
from ..discover import discover_channel
|
||||
from ..discover import discover_incremental
|
||||
from ..pipeline import process_video
|
||||
from ..ratelimit import polite_sleep
|
||||
from ..store import Store
|
||||
from ..ratelimit import ThrottleGuard, polite_sleep
|
||||
from ..store import Store, VideoRef
|
||||
|
||||
|
||||
class JobManager:
|
||||
@@ -89,12 +90,18 @@ class JobManager:
|
||||
if not job:
|
||||
return
|
||||
opts = json.loads(job.opts_json) if job.opts_json else {}
|
||||
if opts.get("video_ids"):
|
||||
self._run_batch(job_id, opts)
|
||||
if opts.get("mode") == "discover":
|
||||
self._run_discovery(job_id, opts)
|
||||
return
|
||||
if opts.get("mode") == "audio":
|
||||
self._run_audio(job_id, opts)
|
||||
return
|
||||
if opts.get("mode") == "video":
|
||||
self._run_video(job_id, opts)
|
||||
return
|
||||
if opts.get("video_ids"):
|
||||
self._run_batch(job_id, opts)
|
||||
return
|
||||
self._run_channel(job_id, opts)
|
||||
|
||||
def _resolve_channel_for_video(self, video_id: str) -> tuple[str, str, str]:
|
||||
@@ -107,19 +114,85 @@ class JobManager:
|
||||
url = f"https://www.youtube.com/@{handle}/videos" if handle else row.url
|
||||
return (name, row.channel_id, url)
|
||||
|
||||
# ---- throttling ------------------------------------------------------
|
||||
#
|
||||
# Every loop that touches YouTube shares these three helpers. Before them,
|
||||
# a throttled session was invisible to the runner: it kept going and turned
|
||||
# one rate-limit into hundreds of videos marked failed (the incident that
|
||||
# left 511 rows in `no_subtitles` and 347 in `error`, all of them saying
|
||||
# "rate-limited by YouTube"). Now the run stops and everything it has not
|
||||
# reached stays `pending`, which is the retryable state.
|
||||
|
||||
def _new_guard(self) -> ThrottleGuard:
|
||||
d = self.cfg.delay
|
||||
return ThrottleGuard(
|
||||
threshold=d.throttle_threshold,
|
||||
base=d.backoff_base,
|
||||
cap=d.backoff_cap,
|
||||
)
|
||||
|
||||
def _note_outcome(self, job_id: str, guard: ThrottleGuard, video_id: str, status: str) -> bool:
|
||||
"""Feed one video's result to the breaker. Returns False to stop the run.
|
||||
|
||||
`process_video` returns only a status string, so the reason is read back
|
||||
from the row: the message it stored is the one place the throttling
|
||||
signature survives.
|
||||
"""
|
||||
if status == "done":
|
||||
guard.note_success()
|
||||
return True
|
||||
row = self.store.get_video(video_id)
|
||||
wait = guard.note_failure(getattr(row, "error_msg", None) or status)
|
||||
if guard.tripped:
|
||||
return False
|
||||
if wait > 0:
|
||||
self._emit(job_id, "log", {
|
||||
"msg": f"YouTube is throttling us — waiting {wait:.1f}s before the next video"
|
||||
})
|
||||
time.sleep(wait)
|
||||
return True
|
||||
|
||||
def _stop_throttled(
|
||||
self, job_id: str, guard: ThrottleGuard, completed: int, total: int,
|
||||
extra: dict[str, Any] | None = None,
|
||||
) -> None:
|
||||
msg = guard.tripped_reason or "stopped by the throttling circuit breaker"
|
||||
remaining = max(0, total - completed)
|
||||
detail = (
|
||||
f"{msg}. Stopped after {completed}/{total}; the remaining {remaining} "
|
||||
"were left untouched and are still pending."
|
||||
)
|
||||
self.store.update_job(
|
||||
job_id, status="error", completed=completed, last_error=detail, finished=True
|
||||
)
|
||||
self._emit(job_id, "log", {"msg": f"ABORTED — {detail}"})
|
||||
self._emit(job_id, "error", {
|
||||
"message": detail, "throttled": True,
|
||||
"completed": completed, "total": total, **(extra or {}),
|
||||
})
|
||||
|
||||
def _run_batch(self, job_id: str, opts: dict[str, Any]) -> None:
|
||||
from ..pipeline import process_video as _process_video, cache_thumbnail
|
||||
"""Generate the .md for an explicit list of videos — and nothing else.
|
||||
|
||||
Thumbnails are deliberately NOT fetched here. They are already cached by
|
||||
the channel-level paths (add-channel, the Thumbnails tool), and
|
||||
/api/thumbnails/{id} redirects to the CDN for anything missing, so
|
||||
piggybacking them on a .md batch only spent extra requests per video for
|
||||
an image the UI could already display.
|
||||
"""
|
||||
from ..pipeline import process_video as _process_video
|
||||
video_ids = opts.get("video_ids") or []
|
||||
cookie_path = opts.get("cookies_file") or resolve_active_path(self.store)
|
||||
thumb_dir = Path(self.cfg.output_dir_resolved).parent / "thumbnails"
|
||||
total = len(video_ids)
|
||||
self.store.update_job(job_id, status="running", total=total, completed=0)
|
||||
self._emit(job_id, "progress", {"completed": 0, "total": total})
|
||||
self._emit(job_id, "log", {"msg": f"processing {total} videos (.md + thumbnails)"})
|
||||
self._emit(job_id, "log", {"msg": f"generating .md for {total} videos"})
|
||||
cfg = _clone_config(self.cfg)
|
||||
if opts.get("languages"):
|
||||
cfg.languages = opts["languages"]
|
||||
cfg.languages = parse_languages(opts["languages"], cfg.prefer_manual)
|
||||
completed = 0
|
||||
outcomes = {"processed": 0, "no_subtitles": 0, "errors": 0}
|
||||
guard = self._new_guard()
|
||||
for vid in video_ids:
|
||||
if job_id in self._cancel:
|
||||
self.store.update_job(job_id, status="cancelled", completed=completed, finished=True)
|
||||
@@ -128,15 +201,17 @@ class JobManager:
|
||||
row = self.store.get_video(vid)
|
||||
if not row:
|
||||
self._emit(job_id, "log", {"msg": f"skip unknown {vid}"})
|
||||
outcomes["errors"] += 1
|
||||
completed += 1
|
||||
self.store.update_job(job_id, completed=completed)
|
||||
continue
|
||||
# UX optimization: skip re-extracting videos whose .md already exists — just ensure thumbnail cached.
|
||||
# UX optimization: a video whose .md is already on disk costs nothing
|
||||
# to "download" — skip the extraction entirely.
|
||||
md_path = Path(cfg.output_dir_resolved).parent / row.markdown_path if row.markdown_path else None
|
||||
if row.status == "done" and row.markdown_path and md_path and md_path.exists():
|
||||
cache_thumbnail(self.store, vid, thumb_dir)
|
||||
status = "done"
|
||||
self._emit(job_id, "log", {"msg": f"cached {vid} (md already present)"})
|
||||
extracted = False
|
||||
self._emit(job_id, "log", {"msg": f"skip {vid} (.md already present)"})
|
||||
else:
|
||||
channel_name, channel_id, channel_url = self._resolve_channel_for_video(vid)
|
||||
status = _process_video(
|
||||
@@ -144,12 +219,28 @@ class JobManager:
|
||||
cookies_file=cookie_path,
|
||||
on_log=lambda m: self._emit(job_id, "log", {"msg": m}),
|
||||
)
|
||||
cache_thumbnail(self.store, vid, thumb_dir)
|
||||
extracted = True
|
||||
if status == "done":
|
||||
outcomes["processed"] += 1
|
||||
elif status == "no_subtitles":
|
||||
outcomes["no_subtitles"] += 1
|
||||
self._emit(job_id, "log", {"msg": f"{vid}: no transcript; .md not generated"})
|
||||
else:
|
||||
outcomes["errors"] += 1
|
||||
self._emit(job_id, "log", {"msg": f"{vid}: processing failed; .md not generated"})
|
||||
completed += 1
|
||||
self.store.update_job(job_id, completed=completed)
|
||||
self._emit(job_id, "progress", {"completed": completed, "total": total, "video_id": vid, "status": status})
|
||||
if not self._note_outcome(job_id, guard, vid, status):
|
||||
self._stop_throttled(job_id, guard, completed, total, outcomes)
|
||||
return
|
||||
# Only pace when we actually hit YouTube. A batch of already-rendered
|
||||
# videos costs nothing but a disk check, and sleeping through it
|
||||
# would make multi-select feel broken for no benefit.
|
||||
if extracted and completed < total:
|
||||
polite_sleep(cfg.delay.min_seconds, cfg.delay.max_seconds)
|
||||
self.store.update_job(job_id, status="done", completed=completed, finished=True)
|
||||
self._emit(job_id, "done", {"completed": completed, "total": total})
|
||||
self._emit(job_id, "done", {"completed": completed, "total": total, **outcomes, "throttling": guard.summary()})
|
||||
|
||||
def _run_audio(self, job_id: str, opts: dict[str, Any]) -> None:
|
||||
import yt_dlp
|
||||
@@ -165,15 +256,28 @@ class JobManager:
|
||||
total = len(videos)
|
||||
self.store.update_job(job_id, status="running", total=total, completed=0)
|
||||
self._emit(job_id, "log", {"msg": f"downloading {total} audio tracks -> {out_dir}"})
|
||||
cfg = self.cfg
|
||||
ydl_opts = {
|
||||
"format": "bestaudio/best",
|
||||
"outtmpl": str(out_dir / "%(id)s.%(ext)s"),
|
||||
"postprocessors": [{"key": "FFmpegExtractAudio", "preferredcodec": "mp3", "preferredquality": "128"}],
|
||||
"quiet": True, "no_warnings": True, "noprogress": True,
|
||||
# This is the only path that downloads media, so it is also the only
|
||||
# one where yt-dlp's own download-side throttles actually fire.
|
||||
"sleep_interval": cfg.delay.min_seconds,
|
||||
"max_sleep_interval": cfg.delay.max_seconds,
|
||||
"sleep_interval_requests": cfg.yt_dlp.sleep_subrequests,
|
||||
"retries": cfg.yt_dlp.retries,
|
||||
"socket_timeout": 30.0,
|
||||
}
|
||||
if cfg.delay.audio_rate_limit:
|
||||
ydl_opts["ratelimit"] = cfg.delay.audio_rate_limit
|
||||
if cookie_path:
|
||||
ydl_opts["cookiefile"] = cookie_path
|
||||
completed = 0
|
||||
guard = self._new_guard()
|
||||
# One YoutubeDL for the whole batch: keeps the connection pool, the
|
||||
# cookie jar and the resolved player JS alive across videos.
|
||||
with yt_dlp.YoutubeDL(ydl_opts) as ydl:
|
||||
for v in videos:
|
||||
if job_id in self._cancel:
|
||||
@@ -183,13 +287,132 @@ class JobManager:
|
||||
try:
|
||||
ydl.download([v.url])
|
||||
self._emit(job_id, "log", {"msg": f"OK {v.video_id}"})
|
||||
guard.note_success()
|
||||
except Exception as exc:
|
||||
self._emit(job_id, "log", {"msg": f"FAIL {v.video_id}: {exc}"})
|
||||
wait = guard.note_failure(exc)
|
||||
if guard.tripped:
|
||||
self._stop_throttled(job_id, guard, completed, total)
|
||||
return
|
||||
if wait > 0:
|
||||
self._emit(job_id, "log", {
|
||||
"msg": f"YouTube is throttling us — waiting {wait:.1f}s before the next track"
|
||||
})
|
||||
time.sleep(wait)
|
||||
completed += 1
|
||||
self.store.update_job(job_id, completed=completed)
|
||||
self._emit(job_id, "progress", {"completed": completed, "total": total, "video_id": v.video_id})
|
||||
self.store.update_job(job_id, status="done", completed=completed, finished=True)
|
||||
self._emit(job_id, "done", {"completed": completed, "total": total})
|
||||
self._emit(job_id, "done", {"completed": completed, "total": total, "throttling": guard.summary()})
|
||||
|
||||
def _run_video(self, job_id: str, opts: dict[str, Any]) -> None:
|
||||
"""Download one validated video as a Chromium-friendly WebM."""
|
||||
import yt_dlp
|
||||
|
||||
video_id = (opts.get("video_ids") or [None])[0]
|
||||
row = self.store.get_video(video_id) if video_id else None
|
||||
if not row:
|
||||
self._fail_video_job(job_id, "video not found")
|
||||
return
|
||||
|
||||
root = Path(self.cfg.output_dir_resolved).parent
|
||||
md_path = root / row.markdown_path if row.markdown_path else None
|
||||
if row.status != "done" or not row.markdown_path or not md_path or not md_path.exists():
|
||||
self._fail_video_job(job_id, "markdown must be generated and present before downloading video", video_id)
|
||||
return
|
||||
|
||||
out_dir = root / "videos" / video_id
|
||||
out_dir.mkdir(parents=True, exist_ok=True)
|
||||
filename = row.video_filename or _safe_video_filename(row.title or video_id)
|
||||
target = out_dir / filename
|
||||
if target.suffix.lower() != ".webm":
|
||||
target = target.with_suffix(".webm")
|
||||
rel_path = target.relative_to(root).as_posix()
|
||||
|
||||
if target.exists() and target.stat().st_size > 0:
|
||||
self.store.update_video_download(video_id, "done", path=rel_path, filename=target.name, size=target.stat().st_size)
|
||||
self.store.update_job(job_id, status="done", total=1, completed=1, finished=True)
|
||||
self._emit(job_id, "progress", {"completed": 1, "total": 1, "video_id": video_id, "status": "done", "skipped": True})
|
||||
self._emit(job_id, "done", {"completed": 1, "total": 1, "skipped": True})
|
||||
return
|
||||
|
||||
cookie_path = opts.get("cookies_file") or resolve_active_path(self.store)
|
||||
restricted = bool(row.block_reason or row.availability in {
|
||||
"subscriber_only", "premium_only", "private", "needs_auth",
|
||||
})
|
||||
self.store.update_video_download(video_id, "queued", filename=target.name, error=None)
|
||||
self.store.update_job(job_id, status="running", total=1, completed=0)
|
||||
self._emit(job_id, "log", {"msg": f"downloading {video_id} -> {target.name}"})
|
||||
|
||||
def hook(data: dict[str, Any]) -> None:
|
||||
if data.get("status") not in ("downloading", "finished"):
|
||||
return
|
||||
downloaded = int(data.get("downloaded_bytes") or 0)
|
||||
total = int(data.get("total_bytes") or data.get("total_bytes_estimate") or 0)
|
||||
percent = round(downloaded * 100 / total, 1) if total else None
|
||||
self._emit(job_id, "progress", {
|
||||
"completed": 0, "total": 1, "video_id": video_id,
|
||||
"status": data.get("status"), "downloaded_bytes": downloaded,
|
||||
"total_bytes": total, "percent": percent,
|
||||
"speed": data.get("speed"), "eta": data.get("eta"),
|
||||
})
|
||||
|
||||
self.store.update_video_download(video_id, "downloading", filename=target.name, error=None)
|
||||
ydl_opts = {
|
||||
# Prefer YouTube's HLS VP9 + Opus pair. Chromium can play the
|
||||
# resulting WebM directly, and HLS avoids the recurring mid-range
|
||||
# 403s seen on long HTTPS media requests.
|
||||
"format": "bestvideo[height<=1080][protocol^=m3u8_native][vcodec^=vp09]+bestaudio[acodec^=opus]/bestvideo[height<=1080][protocol^=m3u8_native]+bestaudio[acodec^=opus]/bestvideo[height<=1080][ext=webm]+bestaudio[ext=webm]",
|
||||
"merge_output_format": "webm",
|
||||
"outtmpl": str(out_dir / f"{filename.rsplit('.', 1)[0]}.%(ext)s"),
|
||||
"quiet": True, "no_warnings": True, "noprogress": True,
|
||||
"progress_hooks": [hook],
|
||||
"continuedl": True,
|
||||
# Current YouTube extraction requires an external JS runtime and
|
||||
# the EJS challenge scripts. Node is available in the supported
|
||||
# local browser setup and is installed by yt-dlp[default].
|
||||
"js_runtimes": {"node": {}},
|
||||
"fragment_retries": 10,
|
||||
"retries": self.cfg.yt_dlp.retries,
|
||||
"extractor_retries": self.cfg.yt_dlp.extractor_retries,
|
||||
"sleep_interval": self.cfg.delay.min_seconds,
|
||||
"max_sleep_interval": self.cfg.delay.max_seconds,
|
||||
"sleep_interval_requests": self.cfg.yt_dlp.sleep_subrequests,
|
||||
"socket_timeout": self.cfg.yt_dlp.socket_timeout,
|
||||
"extractor_args": {"youtube": {"player_client": ["default", "web"]}},
|
||||
}
|
||||
if cookie_path and restricted:
|
||||
ydl_opts["cookiefile"] = cookie_path
|
||||
|
||||
try:
|
||||
with yt_dlp.YoutubeDL(ydl_opts) as ydl:
|
||||
ydl.download([row.url])
|
||||
if not target.exists():
|
||||
# yt-dlp may retain the requested stem but choose a different
|
||||
# extension when the merge was skipped; locate only this ID's
|
||||
# directory and accept the generated WebM as the canonical file.
|
||||
candidates = sorted(out_dir.glob("*.webm"), key=lambda p: p.stat().st_mtime, reverse=True)
|
||||
if candidates:
|
||||
target = candidates[0]
|
||||
if not target.exists() or target.stat().st_size <= 0:
|
||||
raise RuntimeError("yt-dlp finished without producing a WebM file")
|
||||
rel_path = target.relative_to(root).as_posix()
|
||||
size = target.stat().st_size
|
||||
self.store.update_video_download(video_id, "done", path=rel_path, filename=target.name, size=size, error=None)
|
||||
self.store.update_job(job_id, status="done", total=1, completed=1, finished=True)
|
||||
self._emit(job_id, "progress", {"completed": 1, "total": 1, "video_id": video_id, "status": "done", "percent": 100, "total_bytes": size})
|
||||
self._emit(job_id, "done", {"completed": 1, "total": 1, "video_id": video_id, "size": size})
|
||||
except Exception as exc:
|
||||
msg = str(exc)
|
||||
self.store.update_video_download(video_id, "error", filename=target.name, error=msg)
|
||||
self.store.update_job(job_id, status="error", total=1, completed=0, last_error=msg, finished=True)
|
||||
self._emit(job_id, "error", {"message": msg, "video_id": video_id})
|
||||
|
||||
def _fail_video_job(self, job_id: str, message: str, video_id: str | None = None) -> None:
|
||||
if video_id:
|
||||
self.store.update_video_download(video_id, "error", error=message)
|
||||
self.store.update_job(job_id, status="error", total=1, completed=0, last_error=message, finished=True)
|
||||
self._emit(job_id, "error", {"message": message, **({"video_id": video_id} if video_id else {})})
|
||||
|
||||
def _run_channel(self, job_id: str, opts: dict[str, Any]) -> None:
|
||||
job = self.store.get_job(job_id)
|
||||
@@ -204,37 +427,75 @@ class JobManager:
|
||||
return
|
||||
cfg.channel_url = channel_url
|
||||
|
||||
self.store.update_job(job_id, status="running")
|
||||
self._emit(job_id, "log", {"msg": f"discovering {channel_url}"})
|
||||
try:
|
||||
_channel_id, channel_name, refs = discover_channel(channel_url, sleep_subrequests=cfg.yt_dlp.sleep_subrequests)
|
||||
except Exception as exc:
|
||||
self.store.update_job(job_id, status="error", last_error=str(exc), finished=True)
|
||||
self._emit(job_id, "error", {"message": f"discovery failed: {exc}"})
|
||||
return
|
||||
|
||||
limit = opts.get("limit")
|
||||
since = opts.get("since")
|
||||
languages = opts.get("languages")
|
||||
include_shorts = opts.get("include_shorts", cfg.include_shorts)
|
||||
no_live = opts.get("no_live", not cfg.include_live)
|
||||
cookie_override = opts.get("cookies_file")
|
||||
full = bool(opts.get("full"))
|
||||
|
||||
sync = cfg.sync
|
||||
incremental = sync.incremental and not full and bool(channel_id)
|
||||
known = self.store.known_video_ids(channel_id) if incremental else set()
|
||||
|
||||
self.store.update_job(job_id, status="running")
|
||||
self._emit(job_id, "log", {
|
||||
"msg": f"discovering {channel_url}"
|
||||
+ (f" (incremental — newest {sync.window} first)" if known else " (full)")
|
||||
})
|
||||
try:
|
||||
result = discover_incremental(
|
||||
channel_url,
|
||||
known,
|
||||
sleep_subrequests=cfg.yt_dlp.sleep_subrequests,
|
||||
window=sync.window,
|
||||
max_window=sync.max_window,
|
||||
overlap=sync.overlap,
|
||||
since=self.store.latest_upload_date(channel_id) if incremental else None,
|
||||
keep=_keep_ref(cfg, include_shorts, no_live),
|
||||
)
|
||||
except Exception as exc:
|
||||
self.store.update_job(job_id, status="error", last_error=str(exc), finished=True)
|
||||
self._emit(job_id, "error", {"message": f"discovery failed: {exc}"})
|
||||
return
|
||||
|
||||
_channel_id, channel_name, avatar = result.channel_id, result.channel_name, result.avatar
|
||||
refs = result.refs
|
||||
self._emit(job_id, "log", {
|
||||
"msg": f"fetched {result.fetched} entries in {result.passes} pass(es), {result.new_count} new"
|
||||
})
|
||||
|
||||
if languages:
|
||||
cfg.languages = languages
|
||||
cfg.languages = parse_languages(languages, cfg.prefer_manual)
|
||||
if since:
|
||||
refs = [r for r in refs if (r.upload_date or "") >= since.replace("-", "")]
|
||||
if not include_shorts:
|
||||
refs = [r for r in refs if "/shorts/" not in (r.url or "")]
|
||||
if no_live:
|
||||
refs = [r for r in refs if not (r.url or "").startswith("https://www.youtube.com/live/")]
|
||||
if limit:
|
||||
refs = refs[: int(limit)]
|
||||
# Keep undated entries: flat discovery does not report upload_date,
|
||||
# so `>= since` on a missing date would discard the whole channel.
|
||||
cutoff = since.replace("-", "")
|
||||
refs = [r for r in refs if not r.upload_date or r.upload_date >= cutoff]
|
||||
|
||||
self.store.upsert_channel(_channel_id, _handle(channel_url), channel_name, len(refs))
|
||||
if not avatar:
|
||||
from ..discover import deep_channel_avatar
|
||||
avatar = deep_channel_avatar(channel_url, sleep_subrequests=cfg.yt_dlp.sleep_subrequests)
|
||||
if result.full_scan:
|
||||
self.store.upsert_channel(_channel_id, _handle(channel_url), channel_name, len(refs), avatar=avatar)
|
||||
else:
|
||||
self.store.update_channel_meta(_channel_id, name=channel_name, avatar=avatar)
|
||||
if avatar:
|
||||
from ..pipeline import cache_channel_avatar
|
||||
cache_channel_avatar(_channel_id, avatar, Path(self.cfg.output_dir_resolved).parent / "avatars")
|
||||
self.store.upsert_videos(refs)
|
||||
self.store.mark_channel_synced(_channel_id)
|
||||
|
||||
pending = [r for r in self.store.get_pending(_channel_id) if r.video_id in {x.video_id for x in refs}]
|
||||
# Process the freshly-seen window first, then the older backlog, so a
|
||||
# `limit` still means "the newest N" now that discovery stops early.
|
||||
pending = _order_pending(self.store.get_pending(_channel_id), refs)
|
||||
pending = [r for r in pending if _keep_ref(cfg, include_shorts, no_live)(r)]
|
||||
if since:
|
||||
cutoff = since.replace("-", "")
|
||||
pending = [r for r in pending if not r.upload_date or r.upload_date >= cutoff]
|
||||
if limit:
|
||||
pending = pending[: int(limit)]
|
||||
total = len(pending)
|
||||
self.store.update_job(job_id, total=total)
|
||||
self._emit(job_id, "progress", {"completed": 0, "total": total})
|
||||
@@ -242,6 +503,8 @@ class JobManager:
|
||||
|
||||
cookie_path = cookie_override or resolve_active_path(self.store)
|
||||
completed = 0
|
||||
outcomes = {"processed": 0, "no_subtitles": 0, "errors": 0}
|
||||
guard = self._new_guard()
|
||||
for row in pending:
|
||||
if job_id in self._cancel:
|
||||
self.store.update_job(job_id, status="cancelled", completed=completed, finished=True)
|
||||
@@ -252,14 +515,193 @@ class JobManager:
|
||||
cookies_file=cookie_path,
|
||||
on_log=lambda m: self._emit(job_id, "log", {"msg": m}),
|
||||
)
|
||||
if status == "done":
|
||||
outcomes["processed"] += 1
|
||||
elif status == "no_subtitles":
|
||||
outcomes["no_subtitles"] += 1
|
||||
self._emit(job_id, "log", {"msg": f"{row.video_id}: no transcript; .md not generated"})
|
||||
else:
|
||||
outcomes["errors"] += 1
|
||||
self._emit(job_id, "log", {"msg": f"{row.video_id}: processing failed; .md not generated"})
|
||||
completed += 1
|
||||
self.store.update_job(job_id, completed=completed)
|
||||
self._emit(job_id, "progress", {"completed": completed, "total": total, "video_id": row.video_id, "status": status})
|
||||
if not self._note_outcome(job_id, guard, row.video_id, status):
|
||||
self._stop_throttled(job_id, guard, completed, total, outcomes)
|
||||
return
|
||||
if completed < total:
|
||||
polite_sleep(cfg.delay.min_seconds, cfg.delay.max_seconds)
|
||||
|
||||
self.store.update_job(job_id, status="done", completed=completed, finished=True)
|
||||
self._emit(job_id, "done", {"completed": completed, "total": total})
|
||||
self._emit(job_id, "done", {"completed": completed, "total": total, **outcomes, "throttling": guard.summary()})
|
||||
|
||||
def _run_discovery(self, job_id: str, opts: dict[str, Any] | None = None) -> None:
|
||||
"""Refresh the catalog without extracting or downloading video content."""
|
||||
job = self.store.get_job(job_id)
|
||||
if not job:
|
||||
return
|
||||
opts = opts or {}
|
||||
full = bool(opts.get("full"))
|
||||
|
||||
if job.channel_id:
|
||||
channel = self.store.get_channel(job.channel_id)
|
||||
channels = [channel] if channel else []
|
||||
else:
|
||||
channels = self.store.list_channels()
|
||||
channels = [c for c in channels if c]
|
||||
total = len(channels)
|
||||
self.store.update_job(job_id, status="running", total=total, completed=0)
|
||||
mode = "full rescan" if full else "incremental (recent uploads only)"
|
||||
self._emit(job_id, "log", {"msg": f"investigating {total} channel(s) — {mode}"})
|
||||
|
||||
totals = {"new_videos": 0, "known_videos": 0, "errors": 0, "fetched": 0}
|
||||
completed = 0
|
||||
guard = self._new_guard()
|
||||
for channel in channels:
|
||||
if job_id in self._cancel:
|
||||
self.store.update_job(job_id, status="cancelled", completed=completed, finished=True)
|
||||
self._emit(job_id, "cancelled", {"completed": completed, "total": total})
|
||||
return
|
||||
|
||||
# Pace between channels too: each one is a fresh burst of
|
||||
# continuation requests, and a catalog refresh over every channel
|
||||
# used to fire them back to back.
|
||||
if completed:
|
||||
polite_sleep(self.cfg.delay.min_seconds, self.cfg.delay.max_seconds)
|
||||
|
||||
channel_id = channel["channel_id"]
|
||||
try:
|
||||
result = self._discover_catalog_channel(channel, full=full)
|
||||
except Exception as exc:
|
||||
wait = guard.note_failure(exc)
|
||||
if guard.tripped:
|
||||
self._stop_throttled(job_id, guard, completed, total, totals)
|
||||
return
|
||||
if wait > 0:
|
||||
self._emit(job_id, "log", {
|
||||
"msg": f"YouTube is throttling us — waiting {wait:.1f}s before the next channel"
|
||||
})
|
||||
time.sleep(wait)
|
||||
if job.channel_id:
|
||||
self.store.update_job(job_id, status="error", completed=completed, last_error=str(exc), finished=True)
|
||||
self._emit(job_id, "error", {"message": f"discovery failed: {exc}"})
|
||||
return
|
||||
totals["errors"] += 1
|
||||
self._emit(job_id, "log", {"msg": f"discovery failed for {channel.get('name') or channel_id}: {exc}"})
|
||||
else:
|
||||
guard.note_success()
|
||||
totals["new_videos"] += result["new_videos"]
|
||||
totals["known_videos"] += result["known_videos"]
|
||||
totals["fetched"] += result["fetched"]
|
||||
name = channel.get("name") or channel_id
|
||||
scope = (
|
||||
f"scanned {result['fetched']} newest"
|
||||
if result["incremental"] else f"scanned all {result['fetched']}"
|
||||
)
|
||||
if not result["caught_up"]:
|
||||
scope += " (hit window ceiling — run a full rescan if videos are missing)"
|
||||
self._emit(job_id, "log", {"msg": f"{name}: {scope}, {result['new_videos']} new"})
|
||||
self._emit(
|
||||
job_id,
|
||||
"progress",
|
||||
{
|
||||
"channel_id": channel_id,
|
||||
"new_videos": result["new_videos"],
|
||||
"known_videos": result["known_videos"],
|
||||
"fetched": result["fetched"],
|
||||
"incremental": result["incremental"],
|
||||
"caught_up": result["caught_up"],
|
||||
"last_video_date": result["last_video_date"],
|
||||
"completed": completed + 1,
|
||||
"total": total,
|
||||
},
|
||||
)
|
||||
completed += 1
|
||||
self.store.update_job(job_id, completed=completed)
|
||||
|
||||
self.store.update_job(job_id, status="done", completed=completed, finished=True)
|
||||
self._emit(job_id, "done", {"completed": completed, "total": total, **totals})
|
||||
|
||||
def _discover_catalog_channel(self, channel: dict, *, full: bool = False) -> dict[str, Any]:
|
||||
"""Refresh one channel's catalog.
|
||||
|
||||
Incremental by default: only the newest slice of the channel is fetched,
|
||||
stopping at the first run of videos we already have. `full` forces the
|
||||
old behaviour (walk every page) for when the local catalog is suspect.
|
||||
"""
|
||||
channel_id = channel["channel_id"]
|
||||
channel_url = _resolve_channel_url(self.store, self.cfg, channel_id)
|
||||
if not channel_url:
|
||||
raise RuntimeError("no channel url")
|
||||
|
||||
sync = self.cfg.sync
|
||||
incremental = sync.incremental and not full
|
||||
known = self.store.known_video_ids(channel_id) if incremental else set()
|
||||
|
||||
result = discover_incremental(
|
||||
channel_url,
|
||||
known,
|
||||
sleep_subrequests=self.cfg.yt_dlp.sleep_subrequests,
|
||||
window=sync.window,
|
||||
max_window=sync.max_window,
|
||||
overlap=sync.overlap,
|
||||
since=self.store.latest_upload_date(channel_id) if incremental else None,
|
||||
keep=_keep_ref(self.cfg),
|
||||
)
|
||||
refs = result.refs
|
||||
|
||||
if result.full_scan:
|
||||
self.store.upsert_channel(result.channel_id, _handle(channel_url), result.channel_name, len(refs))
|
||||
new_videos = self.store.upsert_videos(refs)
|
||||
else:
|
||||
self.store.update_channel_meta(result.channel_id, name=result.channel_name)
|
||||
new_videos = self.store.upsert_videos(refs)
|
||||
marks = self.store.mark_channel_synced(result.channel_id)
|
||||
|
||||
return {
|
||||
"new_videos": new_videos,
|
||||
"known_videos": len({r.video_id for r in refs}) - new_videos,
|
||||
"fetched": result.fetched,
|
||||
"passes": result.passes,
|
||||
"incremental": not result.full_scan,
|
||||
"caught_up": result.caught_up,
|
||||
"last_video_date": marks["last_video_date"],
|
||||
}
|
||||
|
||||
|
||||
def _keep_ref(cfg: Config, include_shorts: bool | None = None, no_live: bool | None = None):
|
||||
"""Predicate matching the shorts/live rules that decide what reaches the DB.
|
||||
|
||||
Incremental discovery needs the same filter its stored ids were created
|
||||
under, otherwise the tail of a window is full of entries that can never be
|
||||
recognised as known and the window keeps widening for nothing.
|
||||
"""
|
||||
shorts = cfg.include_shorts if include_shorts is None else include_shorts
|
||||
skip_live = (not cfg.include_live) if no_live is None else no_live
|
||||
|
||||
def keep(r: Any) -> bool: # VideoRef or VideoRow — both carry `.url`
|
||||
url = r.url or ""
|
||||
if not shorts and "/shorts/" in url:
|
||||
return False
|
||||
if skip_live and url.startswith("https://www.youtube.com/live/"):
|
||||
return False
|
||||
return True
|
||||
|
||||
return keep
|
||||
|
||||
|
||||
def _order_pending(rows: list, refs: list[VideoRef]) -> list:
|
||||
"""Newest-window-first ordering for the pending queue.
|
||||
|
||||
`get_pending` is ordered by discovery time, which used to coincide with
|
||||
newest-first because discovery saw the whole channel at once. With windowed
|
||||
sync that no longer holds, so put the videos from this run's window at the
|
||||
front and keep the rest of the backlog behind them.
|
||||
"""
|
||||
by_id = {r.video_id: r for r in rows}
|
||||
ordered = [by_id.pop(x.video_id) for x in refs if x.video_id in by_id]
|
||||
ordered.extend(by_id.values())
|
||||
return ordered
|
||||
|
||||
|
||||
def _clone_config(cfg: Config) -> Config:
|
||||
@@ -283,3 +725,13 @@ def _handle(url: str) -> str:
|
||||
if "@" in url:
|
||||
return "@" + url.split("@", 1)[1].split("/", 1)[0]
|
||||
return ""
|
||||
|
||||
|
||||
def _safe_video_filename(title: str) -> str:
|
||||
"""Keep the displayed title while making a valid, bounded Windows name."""
|
||||
invalid = set(r'\\/:*?"<>|')
|
||||
safe = "".join("_" if c in invalid or ord(c) < 32 else c for c in title)
|
||||
safe = " ".join(safe.strip().split())
|
||||
safe = safe.rstrip(" .") or "video"
|
||||
# Leave room for the id directory and yt-dlp's temporary suffixes.
|
||||
return safe[:180] + ".webm"
|
||||
|
||||
@@ -38,7 +38,7 @@
|
||||
channels: { items: [], pending: {} },
|
||||
videos: { items: [], total: 0, page: 1, size: 25, selected: [] },
|
||||
sidebarOpen: false,
|
||||
detail: { video: null, transcript: [], hasAudio: false, audioPlaying: false, activeSeg: -1 },
|
||||
detail: { video: null, transcript: [], hasAudio: false, audioPlaying: false, activeSeg: -1, videoError: false },
|
||||
search: { q: "", channel: "", items: [], ran: false },
|
||||
analysis: { tab: "wordcloud", channel: "", term: "", words: [], timeline: [] },
|
||||
|
||||
@@ -330,7 +330,7 @@
|
||||
this.view = "detail";
|
||||
this.loading.detail = true;
|
||||
this.loading.transcript = true;
|
||||
this.detail = { video: null, transcript: [], hasAudio: false, audioPlaying: false, activeSeg: -1 };
|
||||
this.detail = { video: null, transcript: [], hasAudio: false, audioPlaying: false, activeSeg: -1, videoError: false };
|
||||
try {
|
||||
const v = await this.api("/api/videos/" + encodeURIComponent(id));
|
||||
// chapters may come embedded or be absent
|
||||
@@ -365,6 +365,34 @@
|
||||
} catch (e) { this.toast("Audio failed: " + e.message, "error"); }
|
||||
},
|
||||
|
||||
async downloadVideoOne(id) {
|
||||
if (!this.detail.video || this.detail.video.video_id !== id) return;
|
||||
if (this.detail.video.status !== "done") {
|
||||
this.toast("Generate the .md before downloading the video", "error");
|
||||
return;
|
||||
}
|
||||
try {
|
||||
const d = await this.api("/api/tools/video", { method: "POST", headers: { "Content-Type": "application/json" }, body: JSON.stringify({ video_ids: [id] }) });
|
||||
this.detail.video.video_download_status = "queued";
|
||||
this.detail.video.video_error = null;
|
||||
this.detail.videoError = false;
|
||||
this.toast("Video download started");
|
||||
this.subscribeJob(d.job_id, "video");
|
||||
} catch (e) { this.toast("Video download failed: " + e.message, "error"); }
|
||||
},
|
||||
|
||||
openVideoPlayer() {
|
||||
this.detail.videoError = false;
|
||||
this.$nextTick(() => {
|
||||
const player = this.$refs.videoPlayer;
|
||||
if (player) { player.load(); player.play().catch(() => {}); }
|
||||
});
|
||||
},
|
||||
|
||||
onVideoError() {
|
||||
this.detail.videoError = true;
|
||||
},
|
||||
|
||||
_pollAudio(id) {
|
||||
if (this._audioPoll) clearInterval(this._audioPoll);
|
||||
let attempts = 0;
|
||||
@@ -928,7 +956,16 @@
|
||||
const es = new EventSource("/api/scrape/" + encodeURIComponent(jobId) + "/stream");
|
||||
this.scrape.es = es;
|
||||
es.addEventListener("log", e => { try { const d = JSON.parse(e.data); this.scrape.log.push(d.msg || ""); } catch (_) {} });
|
||||
es.addEventListener("progress", e => { try { const d = JSON.parse(e.data); this.scrape.progress = d; } catch (_) {} });
|
||||
es.addEventListener("progress", e => {
|
||||
try {
|
||||
const d = JSON.parse(e.data);
|
||||
this.scrape.progress = d;
|
||||
if (kind === "video" && this.detail.video && this.detail.video.video_id === d.video_id) {
|
||||
this.detail.video.video_download_status = d.status === "finished" || d.status === "done" ? "downloading" : "downloading";
|
||||
this.detail.video.video_progress = d;
|
||||
}
|
||||
} catch (_) {}
|
||||
});
|
||||
es.addEventListener("done", e => {
|
||||
let result = {};
|
||||
try { result = e.data ? JSON.parse(e.data) : {}; } catch (_) {}
|
||||
@@ -937,7 +974,9 @@
|
||||
this.scrape.log.push("[done] job finished");
|
||||
this.closeStream();
|
||||
this.loadJobs();
|
||||
if (this.scrape.kind === "audio") {
|
||||
if (this.scrape.kind === "video") {
|
||||
this.toast("Video download complete", "success");
|
||||
} else if (this.scrape.kind === "audio") {
|
||||
this.toast("Audio job complete", "success");
|
||||
} else if (this.scrape.kind === "discovery") {
|
||||
const n = Number(result.new_videos || 0);
|
||||
@@ -1238,6 +1277,19 @@
|
||||
return m.toFixed(0) + "m";
|
||||
},
|
||||
|
||||
fmtBytes(n) {
|
||||
n = Number(n || 0);
|
||||
if (!n) return "—";
|
||||
const units = ["B", "KB", "MB", "GB", "TB"];
|
||||
let i = 0;
|
||||
while (n >= 1024 && i < units.length - 1) { n /= 1024; i++; }
|
||||
return (i ? n.toFixed(n >= 10 ? 0 : 1) : Math.round(n)) + " " + units[i];
|
||||
},
|
||||
|
||||
fmtSpeed(n) {
|
||||
return n ? this.fmtBytes(n) + "/s" : "—";
|
||||
},
|
||||
|
||||
fmtNum(n) {
|
||||
n = Number(n || 0);
|
||||
if (!isFinite(n)) return "0";
|
||||
|
||||
@@ -406,9 +406,9 @@
|
||||
</button>
|
||||
<button class="btn accent-grad btn-primary !py-1.5 !px-3 text-xs" @click="downloadVideoOne(detail.video.video_id)" x-show="detail.video.status==='done' && detail.video.video_download_status!=='done'" :disabled="detail.video.video_download_status==='queued' || detail.video.video_download_status==='downloading'">
|
||||
<svg x-show="detail.video.video_download_status==='queued' || detail.video.video_download_status==='downloading'" class="spin w-3.5 h-3.5" viewBox="0 0 24 24" fill="none"><circle cx="12" cy="12" r="9" stroke="currentColor" stroke-width="3" stroke-dasharray="40 20"/></svg>
|
||||
<span x-text="detail.video.video_download_status==='queued' || detail.video.video_download_status==='downloading' ? 'Downloading video…' : 'Download video · 1080p MKV'"></span>
|
||||
<span x-text="detail.video.video_download_status==='queued' || detail.video.video_download_status==='downloading' ? 'Downloading video…' : 'Download video · 1080p WebM'"></span>
|
||||
</button>
|
||||
<a class="btn btn-ghost !py-1.5 !px-3 text-xs" x-show="detail.video.video_download_status==='done'" :href="'/api/videos/'+detail.video.video_id+'/media'" :download="detail.video.video_filename || ''">Download MKV</a>
|
||||
<a class="btn btn-ghost !py-1.5 !px-3 text-xs" x-show="detail.video.video_download_status==='done'" :href="'/api/videos/'+detail.video.video_id+'/media'" :download="detail.video.video_filename || ''">Download WebM</a>
|
||||
<button class="btn btn-ghost !py-1.5 !px-3 text-xs" @click="openClip(detail.video.video_id)">Clip</button>
|
||||
<template x-if="detail.video.url"><a class="btn btn-ghost !py-1.5 !px-3 text-xs" :href="detail.video.url" target="_blank" rel="noopener">Open on YouTube ↗</a></template>
|
||||
</div>
|
||||
@@ -429,18 +429,18 @@
|
||||
<div class="cinema-card" x-show="detail.video && detail.video.video_download_status==='done'">
|
||||
<div class="cinema-stage">
|
||||
<video x-ref="videoPlayer" class="cinema-video" controls playsinline preload="metadata"
|
||||
:src="'/api/videos/'+detail.video.video_id+'/media'"
|
||||
:src="'/api/videos/'+detail.video.video_id+'/media?inline=1'"
|
||||
@error="onVideoError()"></video>
|
||||
<div class="cinema-error" x-show="detail.videoError">
|
||||
<div class="text-white font-semibold">Este navegador no pudo reproducir el MKV directamente.</div>
|
||||
<div class="text-white font-semibold">Este navegador no pudo reproducir el WebM directamente.</div>
|
||||
<div class="text-sm text-zinc-400 mt-1">Puedes descargar el archivo y abrirlo con VLC u otro reproductor local.</div>
|
||||
<a class="btn btn-ghost mt-3" :href="'/api/videos/'+detail.video.video_id+'/media'" :download="detail.video.video_filename || ''">Descargar MKV</a>
|
||||
<a class="btn btn-ghost mt-3" :href="'/api/videos/'+detail.video.video_id+'/media'" :download="detail.video.video_filename || ''">Descargar WebM</a>
|
||||
</div>
|
||||
</div>
|
||||
<div class="cinema-meta">
|
||||
<div>
|
||||
<div class="text-white font-semibold" x-text="detail.video.title"></div>
|
||||
<div class="text-xs text-zinc-500 mt-1">1080p · MKV · <span x-text="fmtBytes(detail.video.video_size)"></span></div>
|
||||
<div class="text-xs text-zinc-500 mt-1">1080p · WebM · <span x-text="fmtBytes(detail.video.video_size)"></span></div>
|
||||
</div>
|
||||
<button class="btn btn-ghost !py-1.5 !px-3 text-xs" @click="openVideoPlayer()">Play</button>
|
||||
</div>
|
||||
|
||||
Reference in New Issue
Block a user