diff --git a/src/yt_scraper/discover.py b/src/yt_scraper/discover.py index a229d88..8a40851 100644 --- a/src/yt_scraper/discover.py +++ b/src/yt_scraper/discover.py @@ -1,31 +1,76 @@ from __future__ import annotations import logging -from typing import Any +from dataclasses import dataclass, field +from datetime import datetime, timezone +from typing import Any, Callable, Collection import yt_dlp +from .ratelimit import GLOBAL_PACER, ydl_throttle_opts from .store import VideoRef log = logging.getLogger(__name__) -def discover_channel(channel_url: str, sleep_subrequests: float = 2.0) -> tuple[str, str, list[VideoRef]]: - """Returns (channel_id, channel_name, video_refs).""" +def _entry_upload_date(entry: dict[str, Any]) -> tuple[str | None, int]: + """(YYYYMMDD | None, approx_flag) para un entry flat. + + Con `youtubetab:approximate_date` yt-dlp llena `timestamp` parseando el + texto relativo que YouTube ya manda en el listado ("hace 3 semanas"): la + fecha sale gratis, con la precision del texto (dia para recientes, mas + gruesa para antiguos). Un upload_date crudo siempre gana: es exacto. + """ + exact = entry.get("upload_date") + if exact: + return str(exact), 0 + ts = entry.get("timestamp") + if not ts: + return None, 0 + return datetime.fromtimestamp(int(ts), tz=timezone.utc).strftime("%Y%m%d"), 1 + +# Defaults for incremental sync; `SyncConfig` in config.py is the tunable copy. +DEFAULT_SYNC_WINDOW = 30 +DEFAULT_MAX_WINDOW = 300 +DEFAULT_OVERLAP = 3 + + +def discover_channel( + channel_url: str, + sleep_subrequests: float = 2.0, + limit: int | None = None, +) -> tuple[str, str, str | None, list[VideoRef]]: + """Returns (channel_id, channel_name, avatar_url, video_refs). + + `limit` caps how many entries yt-dlp pulls off the playlist. It is not a + post-filter: yt-dlp stops requesting continuation pages once it has enough, + so a small limit is the difference between one request and dozens. + """ ydl_opts: dict[str, Any] = { "extract_flat": "in_playlist", "quiet": True, "no_warnings": True, "skip_download": True, - "extract_flat_args": None, - "sleep_subrequests": sleep_subrequests, + # Parse the relative time text ("3 weeks ago") YouTube already includes + # in the listing into an approximate `timestamp` per entry. Rides the + # same requests: zero extra calls. + "extractor_args": {"youtubetab": {"approximate_date": ["true"]}}, + **ydl_throttle_opts(sleep_subrequests), } + if limit: + # Only honoured because extract_info processes the result below. With + # `process=False` yt-dlp hands back a lazy generator that ignores + # playlistend, and walking it paginates the entire channel: measured on + # a 2564-video channel, 2 requests versus 86. + ydl_opts["playlistend"] = int(limit) + GLOBAL_PACER.wait() with yt_dlp.YoutubeDL(ydl_opts) as ydl: info = ydl.extract_info(channel_url, download=False) channel_id = info.get("channel_id") or info.get("uploader_id") or info.get("id") or "unknown" channel_name = info.get("channel") or info.get("title") or info.get("uploader") or "Unknown" + avatar = _pick_channel_avatar(info) entries = _flatten_entries(info) refs: list[VideoRef] = [] @@ -36,22 +81,202 @@ def discover_channel(channel_url: str, sleep_subrequests: float = 2.0) -> tuple[ if entry.get("_type") == "playlist": continue url = entry.get("url") or f"https://www.youtube.com/watch?v={video_id}" - upload_date = entry.get("upload_date") + upload_date, date_approx = _entry_upload_date(entry) duration = entry.get("duration") title = entry.get("title") or video_id + # Flat entries report availability ("subscriber_only" for members-only + # content), so a video we could never fetch is identifiable from the + # listing itself instead of costing a failed extraction to discover. + availability = entry.get("availability") refs.append( VideoRef( video_id=str(video_id), channel_id=str(channel_id), title=str(title), url=str(url), - upload_date=str(upload_date) if upload_date else None, + upload_date=upload_date, + date_approx=date_approx, duration=int(duration) if duration else None, + availability=str(availability) if availability else None, + # The tab is reverse-chronological and flat entries carry no + # upload_date, so this index is the only thing that says how old + # an un-extracted video is. Recorded verbatim; later filtering + # (shorts/live, `since`) leaves gaps but never reorders. + position=len(refs), ) ) log.info("Discovered %d videos on channel %s (%s)", len(refs), channel_name, channel_id) - return str(channel_id), str(channel_name), refs + return str(channel_id), str(channel_name), avatar, refs + + +@dataclass +class IncrementalDiscovery: + """Outcome of a windowed channel sync.""" + + channel_id: str + channel_name: str + avatar: str | None + refs: list[VideoRef] = field(default_factory=list) # window fetched, newest first + new_refs: list[VideoRef] = field(default_factory=list) # subset absent from known_ids + fetched: int = 0 # entries yt-dlp actually returned on the last pass + passes: int = 0 # flat extractions performed + window: int = 0 # playlistend used on the last pass + caught_up: bool = False # reached videos we already had (or the channel's end) + exhausted: bool = False # the window covered the entire channel + full_scan: bool = False # we deliberately walked everything + + @property + def new_count(self) -> int: + return len(self.new_refs) + + +def discover_incremental( + channel_url: str, + known_ids: Collection[str], + *, + sleep_subrequests: float = 2.0, + window: int = DEFAULT_SYNC_WINDOW, + max_window: int = DEFAULT_MAX_WINDOW, + overlap: int = DEFAULT_OVERLAP, + since: str | None = None, + keep: Callable[[VideoRef], bool] | None = None, +) -> IncrementalDiscovery: + """Fetch only the newest slice of a channel instead of paginating all of it. + + A channel's /videos tab is reverse-chronological, so once `overlap` + consecutive entries are ones we already have, everything older is already in + the DB and further pages buy nothing but requests against YouTube's limits. + + `since` (YYYYMMDD) is a second, best-effort stop condition for the case where + yt-dlp does report upload dates — the YouTube tab extractor usually does not + populate them in flat mode, which is why the id overlap is what actually + terminates the walk. + + `keep` filters each window before the overlap test, so it must match whatever + filter decided which videos reached the DB (shorts/live); otherwise the tail + would be full of entries that could never be "known" and the window would + widen pointlessly. + """ + known = set(known_ids) + overlap = max(1, int(overlap)) + window = max(1, int(window)) + max_window = max(window, int(max_window)) + + def _apply(raw: list[VideoRef]) -> list[VideoRef]: + return [r for r in raw if keep(r)] if keep else list(raw) + + # No prior state means there is no boundary to find — walk the whole channel. + if not known: + cid, name, avatar, raw = discover_channel(channel_url, sleep_subrequests=sleep_subrequests) + refs = _apply(raw) + log.info("Full discovery of %s: %d videos", channel_url, len(refs)) + return IncrementalDiscovery( + channel_id=cid, channel_name=name, avatar=avatar, + refs=refs, new_refs=list(refs), fetched=len(raw), passes=1, + window=len(raw), caught_up=True, exhausted=True, full_scan=True, + ) + + size = window + passes = 0 + while True: + cid, name, avatar, raw = discover_channel( + channel_url, sleep_subrequests=sleep_subrequests, limit=size + ) + passes += 1 + exhausted = len(raw) < size + refs = _apply(raw) + tail = refs[-overlap:] + + caught_up = exhausted or (bool(tail) and all(r.video_id in known for r in tail)) + if not caught_up and since and tail: + dated = [r.upload_date for r in tail if r.upload_date] + caught_up = bool(dated) and all(d < since for d in dated) + + if caught_up or size >= max_window: + new_refs = [r for r in refs if r.video_id not in known] + log.info( + "Incremental sync of %s: %d fetched over %d pass(es), %d new%s", + channel_url, len(raw), passes, len(new_refs), + "" if caught_up else " (window ceiling hit; older videos not checked)", + ) + return IncrementalDiscovery( + channel_id=cid, channel_name=name, avatar=avatar, + refs=refs, new_refs=new_refs, fetched=len(raw), passes=passes, + window=size, caught_up=caught_up, exhausted=exhausted, + ) + + # Everything in the window was new: the channel published more than we + # looked at, so widen and try again rather than miss uploads. + size = min(size * 2, max_window) + + +def _pick_channel_avatar(info: dict[str, Any]) -> str | None: + """Best-effort channel avatar URL from a yt-dlp channel info dict. + + Only sources that ACTUALLY point to an image are considered. Notably we do + NOT fall back to `channel_url` or `thumbnail` (those are page URLs / banners). + """ + ths = info.get("thumbnails") or [] + if isinstance(ths, list): + for t in reversed(ths): + url = (t.get("url") if isinstance(t, dict) else None) + if isinstance(url, str) and url: + return url + for key in ("avatar", "channel_icon"): + v = info.get(key) + if isinstance(v, str) and v.startswith("http"): + return v + if isinstance(v, dict): + u = v.get("url") + if isinstance(u, str) and u: + return u + return None + + +_CHANNEL_PATH_TAILS = ( + "/videos", "/shorts", "/streams", "/featured", "/playlists", + "/community", "/about", "/channels", +) + + +def _channel_root_url(channel_url: str) -> str: + """Strip a tab suffix from a YouTube channel URL so yt-dlp extracts the + channel home (where the avatar reliably lives), not a tab.""" + u = (channel_url or "").rstrip("/") + for tail in _CHANNEL_PATH_TAILS: + if u.endswith(tail): + u = u[: -len(tail)] + break + return u or channel_url + + +def deep_channel_avatar(channel_url: str, sleep_subrequests: float = 2.0) -> str | None: + """Robust avatar recovery via a yt-dlp call on the channel root URL. + + `playlistend` and `extract_flat` are load-bearing, not tuning. Without them + this asked yt-dlp to fully extract every video the channel has ever + published in order to read one image URL: yt-dlp redirects a bare channel + URL back to /videos, then walks /videos, /streams and /shorts, and + `download=False` suppresses only the media download, not the extraction. + Measured against a 2564-video channel it was still going at 735 requests + when the measurement aborted it; the form below costs 4. + """ + ydl_opts: dict[str, Any] = { + "quiet": True, "no_warnings": True, "skip_download": True, + "extract_flat": "in_playlist", + "playlistend": 1, + **ydl_throttle_opts(sleep_subrequests), + } + target = _channel_root_url(channel_url) + # yt-dlp probes /videos, /streams and /shorts on a bare channel URL. + GLOBAL_PACER.wait(cost=3) + try: + with yt_dlp.YoutubeDL(ydl_opts) as ydl: + info = ydl.extract_info(target, download=False) + except Exception: + return None + return _pick_channel_avatar(info or {}) def _flatten_entries(info: dict[str, Any]) -> list[dict[str, Any]]: diff --git a/src/yt_scraper/store.py b/src/yt_scraper/store.py index 78755b6..c801b99 100644 --- a/src/yt_scraper/store.py +++ b/src/yt_scraper/store.py @@ -49,9 +49,42 @@ _VIDEO_COLUMNS: dict[str, str] = { "description": "TEXT", "chapters_json": "TEXT", "segments_json": "TEXT", + "availability": "TEXT", + # Position in the channel's reverse-chronological /videos tab (higher = + # newer). The only recency signal discovery produces: yt-dlp's flat listing + # reports no upload_date for YouTube entries, so a video that has never been + # extracted has no date to sort by. + "channel_seq": "INTEGER", + # 1 cuando upload_date viene del discovery aproximado (texto relativo de + # YouTube: "hace 3 semanas"), 0/NULL cuando es exacto (extraccion). + "upload_date_approx": "INTEGER DEFAULT 0", +} + +_CHANNEL_COLUMNS: dict[str, str] = { + "avatar": "TEXT", + # Incremental-sync watermark: when we last looked, and the newest upload + # date we know of. `last_video_date` is the "desde aqui en adelante" mark. + "last_synced_at": "TEXT", + "last_video_date": "TEXT", } _EXTRA_SCHEMA = """ +CREATE INDEX IF NOT EXISTS idx_videos_channel_seq ON videos(channel_id, channel_seq); + +-- SORT_DATE_SQL runs two per-channel lookups for every undated row, and undated +-- is the majority of a library until it is fully scraped (4541 of 4959 rows on +-- the real one). Both are PARTIAL and covering, indexing only the dated rows — +-- which is what the lookups are hunting for. Without the partial predicate the +-- "nearest dated video above me" search walks every row in between checking the +-- table for a date it will not find: 2500 rows deep into a channel whose 8 +-- dated videos all sit at the top, that is quadratic, and it measured 576 ms +-- per page against 8 ms with these. Keep the WHERE clauses spelled exactly as +-- the queries spell them or SQLite will not consider the index. +CREATE INDEX IF NOT EXISTS idx_videos_dated_seq ON videos(channel_id, channel_seq, upload_date) + WHERE upload_date IS NOT NULL AND upload_date <> ''; +CREATE INDEX IF NOT EXISTS idx_videos_dated ON videos(channel_id, upload_date) + WHERE upload_date IS NOT NULL AND upload_date <> ''; + CREATE TABLE IF NOT EXISTS transcript_segments ( video_id TEXT NOT NULL, idx INTEGER NOT NULL, @@ -92,6 +125,131 @@ CREATE TABLE IF NOT EXISTS scrape_jobs ( """ +# yt-dlp's availability enum (see yt_dlp/extractor/common.py:414): +# 'private' | 'premium_only' | 'subscriber_only' | 'needs_auth' | 'unlisted' | 'public'. +# Only these four mean "we cannot fetch it" — `unlisted` downloads perfectly +# well and must NOT be treated as blocked. +BLOCKING_AVAILABILITY = { + "subscriber_only": "members_only", + "premium_only": "premium_only", + "private": "private", + "needs_auth": "needs_auth", +} + + +# Sorts below every real YYYYMMDD. Reached only when a video has no date and +# neither does anything else in its channel, i.e. we have zero evidence about +# when it was published. Such a video does not get to outrank videos we do know +# something about; within its channel `channel_seq` still orders it correctly. +# Stripped before it reaches the UI — it is a rank, not a date. +NO_DATE_SENTINEL = "00000000" + +# The date a video is ordered by. +# +# Chronological order here has to mean what it means on YouTube: newest upload +# first, whether or not we have scraped the video. The obstacle is that +# discovery cannot supply `upload_date` — yt-dlp's flat listing does not report +# one for YouTube entries — so every video without a .md also has a NULL date. +# Measured on the real library: 4541 of 4959 rows. +# +# `channel_seq` is the video's position in the channel's reverse-chronological +# /videos tab, which gives an exact within-channel order and a defensible date: +# +# 1. its own upload_date, once an extraction has learned it; +# 2. else the date of the nearest video ABOVE it in the channel that has one — +# it was published no earlier than that, and ties break by rank, so it +# lands in the slot YouTube would give it; +# 3. else the newest date known anywhere in its channel. This is the run at +# the very top of a channel, above every dated video: it is newer than all +# of them (rank settles that) but claiming more would be inventing a date, +# and it used to let a wholly un-scraped channel take over page one; +# 4. else nothing is known at all — see NO_DATE_SENTINEL. +# +# The fallback this replaced was `discovered_at`, which dated every un-scraped +# video "today" and pinned the entire backlog above everything else. +# +# The `upload_date IS NOT NULL AND upload_date <> ''` spelling is load-bearing: +# it is what makes the partial indexes above applicable. +SORT_DATE_SQL = f"""COALESCE( + NULLIF(videos.upload_date, ''), + (SELECT v2.upload_date FROM videos v2 + WHERE v2.channel_id = videos.channel_id + AND v2.channel_seq > COALESCE(videos.channel_seq, -1) + AND v2.upload_date IS NOT NULL AND v2.upload_date <> '' + ORDER BY v2.channel_seq ASC LIMIT 1), + (SELECT MAX(v3.upload_date) FROM videos v3 + WHERE v3.channel_id = videos.channel_id + AND v3.upload_date IS NOT NULL AND v3.upload_date <> ''), + '{NO_DATE_SENTINEL}')""" + +# Ordering references the `sort_date` alias rather than repeating the subquery, +# so every query that uses `_order_clause` must select `SORT_DATE_SQL AS +# sort_date`. +# +# The tiebreak is `channel_id` THEN `channel_seq`, in that order. `channel_seq` +# is a per-channel counter whose maximum is that channel's video count, so +# comparing it ACROSS channels just ranks by catalogue size: measured on the +# real library, 92% of adjacent pairs tie on sort_date, and in the 20250419 tie +# the whole of one channel preceded the whole of another purely because 575 > +# 476. Grouping by channel first keeps each channel's block contiguous and its +# internal order — the part that has to match YouTube — untouched. `video_id` +# then makes the order total, without which LIMIT/OFFSET paging can repeat or +# skip rows between pages. +_NEWEST_FIRST = ( + "sort_date DESC, videos.channel_id, videos.channel_seq DESC, videos.video_id DESC" +) + +# Ascending needs the unknown-date rows pushed out explicitly. Descending gets +# it for free — NO_DATE_SENTINEL sorts below every real date — but that is the +# same reason it sorts FIRST under ASC, which made "oldest" open with the 212 +# videos of a channel nothing has ever extracted, ahead of a genuine 2017 upload. +# "We do not know" is not "the beginning of time"; it belongs at the end either way. +_OLDEST_FIRST = ( + f"(sort_date = '{NO_DATE_SENTINEL}'), " + "sort_date ASC, videos.channel_id, videos.channel_seq ASC, videos.video_id ASC" +) + + +def _nulls_last(column: str, direction: str) -> str: + """`ORDER BY` fragment that keeps NULLs at the bottom either way. + + SQLite only accepts NULLS LAST from 3.30; the boolean-first form works on + every version, and a video with no view count should not outrank one with a + known count just because the column is empty. + """ + return f"videos.{column} IS NULL, videos.{column} {direction}, {_NEWEST_FIRST}" + + +def _order_clause(sort: str | None) -> str: + """Map an API sort key to SQL. Unknown keys fall back to newest-first. + + Both the bare and suffixed spellings are accepted because the web UI sends + `upload_date` / `view_count` while the CLI and older callers send + `upload_date_desc` / `views_desc`; the mismatch used to drop every non-date + sort onto a raw `upload_date DESC` that ignored the inference above. + """ + by_views = _nulls_last("view_count", "DESC") + by_likes = _nulls_last("like_count", "DESC") + return { + "upload_date": _NEWEST_FIRST, + "upload_date_desc": _NEWEST_FIRST, + "newest": _NEWEST_FIRST, + "upload_date_asc": _OLDEST_FIRST, + "oldest": _OLDEST_FIRST, + "duration": _nulls_last("duration", "DESC"), + "duration_desc": _nulls_last("duration", "DESC"), + "duration_asc": _nulls_last("duration", "ASC"), + "view_count": by_views, + "view_count_desc": by_views, + "views_desc": by_views, + "like_count": by_likes, + "like_count_desc": by_likes, + "likes_desc": by_likes, + "title": f"videos.title COLLATE NOCASE ASC, {_NEWEST_FIRST}", + "title_asc": f"videos.title COLLATE NOCASE ASC, {_NEWEST_FIRST}", + }.get((sort or "").strip(), _NEWEST_FIRST) + + @dataclass class VideoRef: video_id: str @@ -100,6 +258,17 @@ class VideoRef: url: str upload_date: str | None = None duration: int | None = None + # Reported by flat discovery, so a members-only video is known before we + # ever spend an extraction attempt on it. + availability: str | None = None + # 0-based index in the listing this ref came from (0 = newest). Discovery + # walks the /videos tab in reverse-chronological order, so this is the + # chronological rank of a video we have no upload_date for yet. + position: int | None = None + # 1 si upload_date es aproximado (derivado del texto relativo del listado), + # 0 si es exacto o no hay fecha. Al final: los llamadores posicionales + # existentes terminan en (upload_date, duration) y no deben desplazarse. + date_approx: int = 0 @dataclass @@ -123,6 +292,35 @@ class VideoRow: description: str | None = None chapters_json: str | None = None segments_json: str | None = None + availability: str | None = None + channel_seq: int | None = None + # 1 cuando upload_date es aproximado (discovery), 0/NULL si es exacto. + upload_date_approx: int | None = None + # Date the row was ordered by. Equals `upload_date` when it is known; for a + # video discovery has not extracted yet it is inferred from `channel_seq` + # (see `query_videos`). Only populated by queries that compute it. + sort_date: str | None = None + + @property + def block_reason(self) -> str | None: + """Why this video can never be fetched, or None if it can. + + Prefers the discovery signal (known before any attempt) and falls back + to the recorded error for rows burned in before availability existed. + """ + blocked = BLOCKING_AVAILABILITY.get((self.availability or "").lower()) + if blocked: + return blocked + msg = (self.error_msg or "").lower() + if not msg: + return None + if "members-only" in msg or "join this channel to get access" in msg: + return "members_only" + if "private video" in msg: + return "private" + if "has been removed" in msg or "has been terminated" in msg: + return "removed" + return None @dataclass @@ -192,7 +390,15 @@ class Store: for col, coltype in _VIDEO_COLUMNS.items(): if col not in existing: conn.execute(f"ALTER TABLE videos ADD COLUMN {col} {coltype}") + ch_existing = {row["name"] for row in conn.execute("PRAGMA table_info(channels)")} + for col, coltype in _CHANNEL_COLUMNS.items(): + if col not in ch_existing: + conn.execute(f"ALTER TABLE channels ADD COLUMN {col} {coltype}") conn.executescript(_EXTRA_SCHEMA) + # Unconditional, not "only when the column was just added": a row can + # also arrive unranked afterwards, and an unranked row is displayed + # in the wrong place rather than merely in an arbitrary one. + _rank_unranked(conn) @contextmanager def _cursor(self) -> Iterator[sqlite3.Cursor]: @@ -205,18 +411,19 @@ class Store: # ------------------------------------------------------------------ channels - def upsert_channel(self, channel_id: str, handle: str | None, name: str | None, video_count: int = 0) -> None: + def upsert_channel(self, channel_id: str, handle: str | None, name: str | None, video_count: int = 0, avatar: str | None = None) -> None: now = _now_iso() with self._cursor() as cur: cur.execute( - """INSERT INTO channels (channel_id, handle, name, last_scraped, video_count) - VALUES (?, ?, ?, ?, ?) + """INSERT INTO channels (channel_id, handle, name, last_scraped, video_count, avatar) + VALUES (?, ?, ?, ?, ?, ?) ON CONFLICT(channel_id) DO UPDATE SET handle = excluded.handle, name = excluded.name, last_scraped = excluded.last_scraped, - video_count = excluded.video_count""", - (channel_id, handle, name, now, video_count), + video_count = excluded.video_count, + avatar = COALESCE(excluded.avatar, channels.avatar)""", + (channel_id, handle, name, now, video_count, avatar), ) def list_channels(self) -> list[dict]: @@ -237,30 +444,166 @@ class Store: cur.execute("DELETE FROM videos WHERE channel_id = ?", (channel_id,)) cur.execute("DELETE FROM channels WHERE channel_id = ?", (channel_id,)) + def known_video_ids(self, channel_id: str) -> set[str]: + """Every video id already recorded for a channel — the boundary an + incremental discovery walks back to.""" + with self._cursor() as cur: + cur.execute("SELECT video_id FROM videos WHERE channel_id = ?", (channel_id,)) + return {row["video_id"] for row in cur.fetchall()} + + def latest_upload_date(self, channel_id: str) -> str | None: + """Newest known upload date (YYYYMMDD) for a channel, or None.""" + with self._cursor() as cur: + cur.execute( + "SELECT MAX(upload_date) AS d FROM videos " + "WHERE channel_id = ? AND upload_date IS NOT NULL AND upload_date <> ''", + (channel_id,), + ) + row = cur.fetchone() + return row["d"] if row and row["d"] else None + + def mark_channel_synced(self, channel_id: str) -> dict[str, Any]: + """Refresh a channel's counters from what is actually stored. + + Incremental discovery only ever sees the newest slice, so `video_count` + has to be recounted here — passing len(refs) would shrink an 848-video + channel to the size of the window. + """ + now = _now_iso() + with self._cursor() as cur: + cur.execute("SELECT COUNT(*) AS n FROM videos WHERE channel_id = ?", (channel_id,)) + count = int(cur.fetchone()["n"]) + cur.execute( + "SELECT MAX(upload_date) AS d FROM videos " + "WHERE channel_id = ? AND upload_date IS NOT NULL AND upload_date <> ''", + (channel_id,), + ) + row = cur.fetchone() + last_date = row["d"] if row and row["d"] else None + cur.execute( + """UPDATE channels + SET video_count = ?, last_scraped = ?, last_synced_at = ?, + last_video_date = COALESCE(?, last_video_date) + WHERE channel_id = ?""", + (count, now, now, last_date, channel_id), + ) + return {"video_count": count, "last_video_date": last_date, "last_synced_at": now} + + def update_channel_meta(self, channel_id: str, *, name: str | None = None, avatar: str | None = None) -> bool: + """Update only the fields explicitly passed. Preserves last_scraped and video_count.""" + sets: list[str] = [] + params: list[Any] = [] + if name is not None: + sets.append("name = ?"); params.append(name) + if avatar is not None: + sets.append("avatar = ?"); params.append(avatar) + if not sets: + return False + params.append(channel_id) + with self._cursor() as cur: + cur.execute(f"UPDATE channels SET {', '.join(sets)} WHERE channel_id = ?", params) + return cur.rowcount > 0 + # ------------------------------------------------------------------ videos def upsert_videos(self, refs: list[VideoRef]) -> int: + """Insert/refresh discovered videos and re-rank the channel's recency order. + + `refs` arrive in /videos-tab order (newest first). That order is the only + chronological signal discovery yields — the flat listing carries no + upload_date — so it is recorded as `channel_seq` and is what lets the UI + place a video that has never been extracted where YouTube would show it. + + The window is ranked ABOVE the channel's current maximum rather than from + zero: a sync only fetches the newest slice, and everything it did not + fetch is by construction older than everything it did. Lifting the window + keeps both halves consistently ordered without re-walking the channel. + """ now = _now_iso() inserted = 0 with self._cursor() as cur: + incoming = {r.video_id for r in refs} + existing: set[str] = set() + if incoming: + placeholders = ",".join("?" for _ in incoming) + cur.execute( + f"SELECT video_id FROM videos WHERE video_id IN ({placeholders})", + list(incoming), + ) + existing = {row["video_id"] for row in cur.fetchall()} + + by_channel: dict[str, list[VideoRef]] = {} + for r in refs: + by_channel.setdefault(r.channel_id, []).append(r) + seqs: dict[str, int] = {} + for channel_id, group in by_channel.items(): + # Honour an explicit position when discovery set one; otherwise + # the list order is the listing order. + ordered = sorted( + enumerate(group), + key=lambda pair: pair[1].position if pair[1].position is not None else pair[0], + ) + cur.execute( + "SELECT COALESCE(MAX(channel_seq), 0) AS m FROM videos WHERE channel_id = ?", + (channel_id,), + ) + base = int(cur.fetchone()["m"] or 0) + width = len(ordered) + for rank, (_, r) in enumerate(ordered): + seqs[r.video_id] = base + width - rank + for r in refs: cur.execute( """INSERT INTO videos - (video_id, channel_id, title, url, upload_date, duration, status, discovered_at) - VALUES (?, ?, ?, ?, ?, ?, 'pending', ?) + (video_id, channel_id, title, url, upload_date, + upload_date_approx, duration, availability, + channel_seq, status, discovered_at) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, 'pending', ?) ON CONFLICT(video_id) DO UPDATE SET - title = excluded.title, - upload_date = excluded.upload_date, - duration = excluded.duration""", - (r.video_id, r.channel_id, r.title, r.url, r.upload_date, r.duration, now), + title = COALESCE(excluded.title, videos.title), + upload_date = CASE + WHEN excluded.upload_date IS NULL THEN videos.upload_date + WHEN videos.upload_date IS NULL THEN excluded.upload_date + WHEN COALESCE(videos.upload_date_approx, 0) = 1 + THEN excluded.upload_date + ELSE videos.upload_date END, + upload_date_approx = CASE + WHEN videos.upload_date IS NOT NULL + AND COALESCE(videos.upload_date_approx, 0) = 0 + THEN videos.upload_date_approx + ELSE COALESCE(excluded.upload_date_approx, + videos.upload_date_approx) END, + duration = COALESCE(excluded.duration, videos.duration), + availability = COALESCE(excluded.availability, videos.availability), + channel_seq = COALESCE(excluded.channel_seq, videos.channel_seq)""", + (r.video_id, r.channel_id, r.title, r.url, r.upload_date, + int(r.date_approx or 0), r.duration, + r.availability, seqs.get(r.video_id), now), ) - if cur.rowcount > 0: + if r.video_id not in existing: inserted += 1 + existing.add(r.video_id) return inserted - def get_pending(self, channel_id: str | None = None, limit: int | None = None) -> list[VideoRow]: + def get_pending( + self, + channel_id: str | None = None, + limit: int | None = None, + *, + include_blocked: bool = False, + ) -> list[VideoRow]: + """Pending videos, excluding ones discovery already told us we cannot + fetch (members-only, premium, private). Bulk runs should not spend + requests on those; an explicit per-video "Process" click still can, + which is what makes the membership case recoverable. + """ sql = "SELECT * FROM videos WHERE status = 'pending'" params: list[Any] = [] + if not include_blocked: + blocking = sorted(BLOCKING_AVAILABILITY) + marks = ",".join("?" for _ in blocking) + sql += f" AND (availability IS NULL OR availability NOT IN ({marks}))" + params.extend(blocking) if channel_id: sql += " AND channel_id = ?" params.append(channel_id) @@ -300,6 +643,7 @@ class Store: sort: str = "upload_date_desc", page: int = 1, size: int = 50, + blocked: bool | None = None, ) -> tuple[list[VideoRow], int]: where: list[str] = [] params: list[Any] = [] @@ -307,6 +651,21 @@ class Store: where.append("channel_id = ?"); params.append(channel_id) if status: where.append("status = ?"); params.append(status) + if blocked is not None: + # Match on the discovery signal or on the recorded error, so rows + # burned in before `availability` existed are still findable. + marks = ",".join("?" for _ in sorted(BLOCKING_AVAILABILITY)) + # COALESCE is load-bearing: with a NULL availability the IN test is + # NULL, and NOT(NULL) is NULL, so the negated branch would silently + # return zero rows instead of "everything fetchable". + expr = ( + f"(COALESCE(availability, '') IN ({marks}) " + "OR COALESCE(error_msg, '') LIKE '%members-only%' " + "OR COALESCE(error_msg, '') LIKE '%Join this channel to get access%' " + "OR COALESCE(error_msg, '') LIKE '%Private video%')" + ) + where.append(expr if blocked else f"NOT {expr}") + params.extend(sorted(BLOCKING_AVAILABILITY)) if date_from: where.append("upload_date >= ?"); params.append(date_from.replace("-", "")) if date_to: @@ -317,19 +676,14 @@ class Store: where.append("(title LIKE ? OR description LIKE ?)") params.extend([f"%{q}%", f"%{q}%"]) clause = ("WHERE " + " AND ".join(where)) if where else "" - order = { - "upload_date_desc": "upload_date DESC", - "upload_date_asc": "upload_date ASC", - "duration_desc": "duration DESC", - "views_desc": "view_count DESC", - "title_asc": "title ASC", - }.get(sort, "upload_date DESC") + order = _order_clause(sort) offset = max(0, (page - 1) * size) with self._cursor() as cur: cur.execute(f"SELECT COUNT(*) AS n FROM videos {clause}", params) total = cur.fetchone()["n"] cur.execute( - f"SELECT * FROM videos {clause} ORDER BY {order} LIMIT ? OFFSET ?", + f"SELECT videos.*, {SORT_DATE_SQL} AS sort_date FROM videos {clause} " + f"ORDER BY {order} LIMIT ? OFFSET ?", [*params, size, offset], ) rows = [_row_to_videorow(r) for r in cur.fetchall()] @@ -392,20 +746,45 @@ class Store: (markdown_path, transcript_lang, transcript_src, int(has_chapters), now, video_id), ) - def mark_status(self, video_id: str, status: str) -> None: + def mark_status(self, video_id: str, status: str, reason: str | None = None) -> None: + """Set a video's status, optionally recording why. + + `no_subtitles` used to be stored with error_msg=NULL, which left no way + to tell "this video has no captions" apart from "the language policy + rejected the captions it does have" — the second is recoverable. + """ now = _now_iso() with self._cursor() as cur: cur.execute( - "UPDATE videos SET status = ?, processed_at = ? WHERE video_id = ?", - (status, now, video_id), + "UPDATE videos SET status = ?, error_msg = ?, processed_at = ? WHERE video_id = ?", + # 2000, not 500: a throttled caption fetch stores the timedtext + # URL, and at 500 the cut landed twenty characters before the + # `tlang=` parameter — the one token proving the request was for + # a machine translation rather than the real transcript. + (status, reason[:2000] if reason else None, now, video_id), + ) + + def set_availability(self, video_id: str, availability: str | None) -> None: + """Record what a full extraction learned; more authoritative than the + flat listing, which omits the field for most entries.""" + if not availability: + return + with self._cursor() as cur: + cur.execute( + "UPDATE videos SET availability = ? WHERE video_id = ?", + (str(availability), video_id), ) def set_upload_date(self, video_id: str, upload_date: str | None) -> None: if not upload_date: return with self._cursor() as cur: + # La fecha de la extraccion es exacta: sobrescribe una aproximada + # del discovery y apaga su bandera; nunca degrada una real. cur.execute( - "UPDATE videos SET upload_date = ? WHERE video_id = ? AND upload_date IS NULL", + """UPDATE videos SET upload_date = ?, upload_date_approx = 0 + WHERE video_id = ? + AND (upload_date IS NULL OR COALESCE(upload_date_approx, 0) = 1)""", (upload_date, video_id), ) @@ -626,27 +1005,6 @@ class Store: r["status"]: r["n"] for r in cur.execute("SELECT status, COUNT(*) AS n FROM videos GROUP BY status").fetchall() } - # uploads over time (by month) - uploads = [ - {"month": r["m"], "count": r["n"]} - for r in cur.execute( - "SELECT substr(upload_date,1,6) AS m, COUNT(*) AS n FROM videos WHERE upload_date IS NOT NULL GROUP BY m ORDER BY m" - ).fetchall() - ] - # duration histogram (buckets) - hist = [ - {"bucket": r["b"], "count": r["n"]} - for r in cur.execute( - """SELECT - CASE WHEN duration < 300 THEN '<5m' - WHEN duration < 600 THEN '5-10m' - WHEN duration < 1200 THEN '10-20m' - WHEN duration < 2400 THEN '20-40m' - ELSE '40m+' END AS b, - COUNT(*) AS n - FROM videos WHERE duration IS NOT NULL GROUP BY b""" - ).fetchall() - ] # top tags (tags is JSON array text) tag_rows = cur.execute("SELECT tags FROM videos WHERE tags IS NOT NULL AND tags != '[]'").fetchall() tag_counts: dict[str, int] = {} @@ -662,14 +1020,59 @@ class Store: return { "channels": channels, "status_breakdown": status_breakdown, - "uploads_over_time": uploads, - "duration_histogram": hist, "top_tags": [{"tag": t, "count": c} for t, c in top_tags], } + RETRYABLE_STATUSES = ("error", "no_subtitles") + + # Failures that will never resolve by trying again: paying for a membership + # or the video coming back from the dead are not retry outcomes. Retrying + # them just spends requests against the rate limit that the videos which + # CAN succeed need. Kept narrow on purpose — YouTube's throttling message + # ("rate-limited ... try again later") is retryable and must not match here. + PERMANENT_ERROR_PATTERNS = ( + "%members-only%", + "%Join this channel to get access%", + "%Private video%", + "%This video has been removed%", + "%video has been terminated%", + ) + + def _permanent_sql(self, negate: bool = True) -> str: + clause = " OR ".join("error_msg LIKE ?" for _ in self.PERMANENT_ERROR_PATTERNS) + return f"NOT (error_msg IS NOT NULL AND ({clause}))" if negate else f"(error_msg IS NOT NULL AND ({clause}))" + def reset_errors(self, channel_id: str | None = None) -> int: - sql = "UPDATE videos SET status = 'pending', error_msg = NULL WHERE status = 'error'" - params: list[Any] = [] + return self.reset_videos(channel_id, statuses=("error",)) + + def reset_videos( + self, + channel_id: str | None = None, + statuses: tuple[str, ...] | list[str] = ("error",), + *, + keep_done: bool = True, + include_permanent: bool = False, + ) -> int: + """Send videos in the given terminal statuses back to `pending`. + + `no_subtitles` has to be resettable, not just `error`: it is recorded + whenever the language policy matched nothing, so a config fix is + worthless if the affected rows can never be retried. `done` is never + reset here — re-running finished work is what burns rate limits. + """ + allowed = [s for s in statuses if s in self.RETRYABLE_STATUSES or not keep_done] + allowed = [s for s in allowed if s != "done"] + if not allowed: + return 0 + placeholders = ",".join("?" for _ in allowed) + sql = ( + f"UPDATE videos SET status = 'pending', error_msg = NULL " + f"WHERE status IN ({placeholders})" + ) + params: list[Any] = list(allowed) + if not include_permanent: + sql += f" AND {self._permanent_sql()}" + params.extend(self.PERMANENT_ERROR_PATTERNS) if channel_id: sql += " AND channel_id = ?" params.append(channel_id) @@ -677,6 +1080,31 @@ class Store: cur.execute(sql, params) return cur.rowcount + def retryable_counts(self, channel_id: str | None = None) -> dict[str, int]: + """Videos stuck in each retryable status, plus how many are permanently + blocked. The retry button must not promise to fix members-only videos. + """ + placeholders = ",".join("?" for _ in self.RETRYABLE_STATUSES) + base = f"FROM videos WHERE status IN ({placeholders})" + base_params: list[Any] = list(self.RETRYABLE_STATUSES) + tail = "" + if channel_id: + tail = " AND channel_id = ?" + with self._cursor() as cur: + cur.execute( + f"SELECT status, COUNT(*) n {base} AND {self._permanent_sql()}{tail} GROUP BY status", + base_params + list(self.PERMANENT_ERROR_PATTERNS) + ([channel_id] if channel_id else []), + ) + out = {s: 0 for s in self.RETRYABLE_STATUSES} + for row in cur.fetchall(): + out[row["status"]] = int(row["n"]) + cur.execute( + f"SELECT COUNT(*) n {base} AND {self._permanent_sql(negate=False)}{tail}", + base_params + list(self.PERMANENT_ERROR_PATTERNS) + ([channel_id] if channel_id else []), + ) + out["permanent"] = int(cur.fetchone()["n"]) + return out + def _sanitize_fts(query: str) -> str: # Build a safe AND FTS5 query from whitespace-separated terms. @@ -695,6 +1123,52 @@ def _now_iso() -> str: return datetime.now(timezone.utc).isoformat(timespec="seconds") +def _rank_unranked(conn: sqlite3.Connection) -> int: + """Give every row that has no recency rank one, above its channel's maximum. + + Covers two cases with the same rule. On a database that predates the column + every row is unranked, and this is the initial seed. Afterwards a row can + still arrive unranked from a process running the pre-`channel_seq` code — a + long-lived server that has not been restarted since the migration — which is + exactly what happened on the live library: three videos discovered after the + migration, all NULL. + + An unranked row is not merely unordered, it is actively misplaced: + `COALESCE(channel_seq, -1)` in SORT_DATE_SQL makes rule 2 pick the OLDEST + dated video in the channel, so a video discovery has only just found — one + of the channel's newest — is shown at the very bottom. Measured: two + brand-new Alex Hormozi uploads displayed with sort_date 20180720, second to + last of 513. + + The ordering is `discovered_at DESC, rowid ASC`, which is the order the rows + were actually learned in: `upsert_videos` stamps one timestamp per discovery + batch, discovery only ever adds ids newer than everything already stored, and + within a batch the insert order is the channel's /videos tab — + reverse-chronological. It is a reconstruction, not an observation; every + later sync overwrites the slice it touches with the real thing. + + Idempotent: with nothing unranked it does no writes at all. + """ + by_channel: dict[str, list[str]] = {} + for row in conn.execute( + "SELECT video_id, channel_id FROM videos WHERE channel_seq IS NULL " + "ORDER BY channel_id, discovered_at DESC, rowid ASC" + ): + by_channel.setdefault(row["channel_id"], []).append(row["video_id"]) + if not by_channel: + return 0 + updates: list[tuple[int, str]] = [] + for channel_id, ids in by_channel.items(): + row = conn.execute( + "SELECT COALESCE(MAX(channel_seq), 0) AS m FROM videos WHERE channel_id = ?", + (channel_id,), + ).fetchone() + base = int(row["m"] or 0) + updates.extend((base + len(ids) - i, vid) for i, vid in enumerate(ids)) + conn.executemany("UPDATE videos SET channel_seq = ? WHERE video_id = ?", updates) + return len(updates) + + def _row_to_videorow(row: sqlite3.Row) -> VideoRow: keys = row.keys() return VideoRow( @@ -717,6 +1191,10 @@ def _row_to_videorow(row: sqlite3.Row) -> VideoRow: description=row["description"] if "description" in keys else None, chapters_json=row["chapters_json"] if "chapters_json" in keys else None, segments_json=row["segments_json"] if "segments_json" in keys else None, + availability=row["availability"] if "availability" in keys else None, + channel_seq=row["channel_seq"] if "channel_seq" in keys else None, + upload_date_approx=row["upload_date_approx"] if "upload_date_approx" in keys else None, + sort_date=row["sort_date"] if "sort_date" in keys else None, ) diff --git a/src/yt_scraper/webapp/api.py b/src/yt_scraper/webapp/api.py index 63543c7..0eab1cb 100644 --- a/src/yt_scraper/webapp/api.py +++ b/src/yt_scraper/webapp/api.py @@ -12,8 +12,15 @@ from .. import analysis as analysis_mod from .. import cookies as cookies_mod from .. import export as export_mod from ..config import Config +from ..ratelimit import polite_sleep +from .. import store as store_mod from ..store import Store +#: How many thumbnails a *implicit* whole-channel request may fetch. The UI +#: pulls the rest lazily through /api/thumbnails/{id} as rows scroll into view, +#: so this only bounds the eager burst that follows "Add channel". +THUMBNAIL_AUTO_LIMIT = 60 + def build_router(store: Store, cfg: Config, jobs) -> APIRouter: r = APIRouter(prefix="/api") @@ -28,16 +35,33 @@ def build_router(store: Store, cfg: Config, jobs) -> APIRouter: def channels_list(): return {"items": store.list_channels()} + # Sync `def`, not `async def`, on purpose: the body blocks on yt-dlp for as + # long as the channel takes to page. FastAPI runs `def` handlers on the + # threadpool, whereas an `async def` would hold the event loop and freeze + # every other request — including the SSE stream of a running job. @r.post("/channels") - async def add_channel(payload: dict): - from ..discover import discover_channel + def add_channel(payload: dict): + from ..discover import discover_channel, deep_channel_avatar + from ..pipeline import cache_channel_avatar url = (payload or {}).get("url") if not url: raise HTTPException(400, "url required") - cid, name, refs = discover_channel(url, sleep_subrequests=cfg.yt_dlp.sleep_subrequests) - store.upsert_channel(cid, _handle(url), name, len(refs)) + cid, name, avatar, refs = discover_channel(url, sleep_subrequests=cfg.yt_dlp.sleep_subrequests) + avatar_cached = False + if not avatar: + avatar = deep_channel_avatar(url, sleep_subrequests=cfg.yt_dlp.sleep_subrequests) + store.upsert_channel(cid, _handle(url), name, len(refs), avatar=avatar) store.upsert_videos(refs) - return {"channel_id": cid, "name": name, "video_count": len(refs)} + if avatar: + avatars_dir = Path(cfg.output_dir_resolved).parent / "avatars" + avatar_cached = cache_channel_avatar(cid, avatar, avatars_dir) + return { + "channel_id": cid, + "name": name, + "video_count": len(refs), + "avatar": avatar, + "avatar_cached": avatar_cached, + } @r.delete("/channels/{channel_id}") def del_channel(channel_id: str): @@ -52,8 +76,11 @@ def build_router(store: Store, cfg: Config, jobs) -> APIRouter: date_to: str | None = Query(None, alias="to"), min_dur: int | None = None, q: str | None = None, sort: str = "upload_date_desc", page: int = 1, size: int = 50, + blocked: bool | None = None, ): - rows, total = store.query_videos(channel, status, date_from, date_to, min_dur, q, sort, page, size) + rows, total = store.query_videos( + channel, status, date_from, date_to, min_dur, q, sort, page, size, blocked=blocked, + ) return {"items": [_video_dict(v) for v in rows], "total": total, "page": page, "size": size} @r.get("/videos/{video_id}") @@ -140,10 +167,18 @@ def build_router(store: Store, cfg: Config, jobs) -> APIRouter: return {"deleted": n} @r.delete("/scrape/{job_id}") - def delete_or_cancel_job(job_id: str): + def delete_or_cancel_job(job_id: str, action: str = "cancel"): job = store.get_job(job_id) if not job: raise HTTPException(404, "job not found") + # action=purge is unconditional: removes the row regardless of status. + # Used by the "Recent jobs" table, where the user's intent is to clear + # history (including rows that survived a server restart as zombies). + # action=cancel (default) only works for terminal jobs by clearing the + # row; for live jobs it asks the JobManager to stop them cooperatively. + if action == "purge": + store.delete_job(job_id) + return {"deleted": job_id} if job.status in store.TERMINAL_STATUSES: store.delete_job(job_id) return {"deleted": job_id} @@ -350,19 +385,31 @@ def build_router(store: Store, cfg: Config, jobs) -> APIRouter: # -------------------------------------------------- thumbnails (local cache) @r.post("/tools/thumbnails") - async def download_thumbnails(payload: dict): + def download_thumbnails(payload: dict): from ..pipeline import cache_thumbnail video_ids = (payload or {}).get("video_ids") or [] channel_id = (payload or {}).get("channel_id") + explicit = bool(video_ids) if channel_id and not video_ids: video_ids = [v.video_id for v in store.get_all(channel_id)] + # Adding a channel fires this from the frontend with just a channel_id, + # which used to mean one CDN request per video in the catalog — ~860 in + # a burst for a large channel, on top of the discovery that just ran. + # Cap the implicit form; an explicit list of ids is the user asking. + limit = int((payload or {}).get("limit") or 0) + if not explicit: + limit = limit or THUMBNAIL_AUTO_LIMIT + skipped = 0 + if limit and len(video_ids) > limit: + skipped = len(video_ids) - limit + video_ids = video_ids[:limit] out_dir = Path(cfg.output_dir_resolved).parent / "thumbnails" out_dir.mkdir(parents=True, exist_ok=True) n = 0 for vid in video_ids: if cache_thumbnail(store, vid, out_dir): n += 1 - return {"downloaded": n, "dir": str(out_dir)} + return {"downloaded": n, "dir": str(out_dir), "skipped": skipped} @r.get("/thumbnails/{video_id}") def serve_thumbnail(video_id: str): @@ -375,6 +422,150 @@ def build_router(store: Store, cfg: Config, jobs) -> APIRouter: url = thumbnail_url_for(v) if v else f"https://i.ytimg.com/vi/{video_id}/hqdefault.jpg" return RedirectResponse(url=url, status_code=302) + # -------------------------------------------------- freshness / recovery + @r.get("/videos-retryable") + def videos_retryable(channel: str | None = None): + """How many videos are stuck in a retryable terminal status. + + `total` excludes the permanently blocked ones (members-only, private, + removed) so the retry button never promises to fix them. + """ + counts = store.retryable_counts(channel) + return { + "counts": counts, + "total": sum(counts.get(s, 0) for s in Store.RETRYABLE_STATUSES), + "permanent": counts.get("permanent", 0), + } + + @r.post("/videos/reset") + async def reset_videos(payload: dict): + """Send `error` / `no_subtitles` videos back to pending so they can be retried. + + `no_subtitles` is resettable on purpose: it is recorded whenever the + language policy matched nothing, so fixing the config is useless if the + affected rows stay terminal. + """ + body = payload or {} + channel_id = body.get("channel_id") + statuses = body.get("statuses") or ["error", "no_subtitles"] + bad = [s for s in statuses if s not in Store.RETRYABLE_STATUSES] + if bad: + raise HTTPException(400, f"not retryable: {bad}") + n = store.reset_videos(channel_id, tuple(statuses)) + return {"reset": n, "statuses": statuses, "channel_id": channel_id} + + @r.post("/tools/reconcile") + def reconcile(payload: dict | None = None): + """Re-scan data/markdown and make the DB agree with disk, on demand. + + Without this the markdown tree is only read at server startup, so any + .md produced afterwards is invisible until a restart. + """ + from ..segments import reconcile_markdown + # prune is opt-in: it demotes `done` rows whose .md vanished, which is + # destructive if the markdown root is ever misconfigured. + prune = bool((payload or {}).get("prune")) + return reconcile_markdown(store, Path(cfg.output_dir_resolved), prune=prune) + + # -------------------------------------------------- channel avatars (local cache) + @r.post("/tools/sync-channels") + def sync_channels(payload: dict): + """Re-scrape each tracked channel to refresh metadata + avatar (with deep fallback).""" + from ..discover import discover_channel, deep_channel_avatar + from ..pipeline import cache_channel_avatar + channel_id = (payload or {}).get("channel_id") + channels = [store.get_channel(channel_id)] if channel_id else store.list_channels() + channels = [c for c in channels if c] + avatars_dir = Path(cfg.output_dir_resolved).parent / "avatars" + avatars_dir.mkdir(parents=True, exist_ok=True) + n_synced = 0 + n_avatars = 0 + n_recovered = 0 + errors: list[dict] = [] + for i, c in enumerate(channels): + # This loop runs outside the JobManager, so nothing else is spacing + # it out; syncing every channel used to be one burst. + if i: + polite_sleep(cfg.delay.min_seconds, cfg.delay.max_seconds) + cid = c["channel_id"] + handle = (c.get("handle") or "").lstrip("@") + ch_url = ( + f"https://www.youtube.com/@{handle}/videos" if handle + else f"https://www.youtube.com/channel/{cid}" + ) + try: + # Metadata + avatar live on the channel object, not the video + # list — one entry is enough and costs a single request instead + # of paginating the entire channel. + _new_id, new_name, avatar, _refs = discover_channel( + ch_url, sleep_subrequests=cfg.yt_dlp.sleep_subrequests, limit=1, + ) + except Exception as exc: + errors.append({"channel_id": cid, "name": c.get("name"), "error": str(exc)}) + continue + if avatar is None: + avatar = deep_channel_avatar(ch_url, sleep_subrequests=cfg.yt_dlp.sleep_subrequests) + if avatar: + n_recovered += 1 + store.update_channel_meta(cid, name=new_name, avatar=avatar) + if avatar and cache_channel_avatar(cid, avatar, avatars_dir): + n_avatars += 1 + n_synced += 1 + return { + "synced": n_synced, + "avatars_cached": n_avatars, + "avatars_recovered_via_deep_fallback": n_recovered, + "errors": errors, + } + + @r.post("/tools/avatars") + def download_avatars(payload: dict): + from ..pipeline import cache_channel_avatar + from ..discover import discover_channel + channel_id = (payload or {}).get("channel_id") + channels = [store.get_channel(channel_id)] if channel_id else store.list_channels() + channels = [c for c in channels if c] + out_dir = Path(cfg.output_dir_resolved).parent / "avatars" + out_dir.mkdir(parents=True, exist_ok=True) + n = 0 + for i, c in enumerate(channels): + if i: + polite_sleep(cfg.delay.min_seconds, cfg.delay.max_seconds) + cid = c["channel_id"] + url = c.get("avatar") + if not url: + handle = (c.get("handle") or "").lstrip("@") + ch_url = f"https://www.youtube.com/@{handle}/videos" if handle else f"https://www.youtube.com/channel/{cid}" + try: + # Avatar only — no reason to walk the video list. + _, _, avatar, _ = discover_channel( + ch_url, sleep_subrequests=cfg.yt_dlp.sleep_subrequests, limit=1, + ) + except Exception: + avatar = None + if avatar: + store.upsert_channel(cid, c.get("handle"), c.get("name"), c.get("video_count") or 0, avatar=avatar) + url = avatar + if url and cache_channel_avatar(cid, url, out_dir): + n += 1 + return {"downloaded": n, "dir": str(out_dir)} + + @r.get("/avatars/{channel_id}") + def serve_avatar(channel_id: str): + from fastapi.responses import RedirectResponse + from ..pipeline import cache_channel_avatar + local = Path(cfg.output_dir_resolved).parent / "avatars" / f"{channel_id}.jpg" + if local.exists(): + return FileResponse(str(local), media_type="image/jpeg") + c = store.get_channel(channel_id) or {} + url = c.get("avatar") + if not url: + raise HTTPException(404, "no avatar for channel") + avatars_dir = Path(cfg.output_dir_resolved).parent / "avatars" + if cache_channel_avatar(channel_id, url, avatars_dir) and local.exists(): + return FileResponse(str(local), media_type="image/jpeg") + return RedirectResponse(url=url, status_code=302) + # -------------------------------------------------- clip (transcript segment) @r.get("/clip/{video_id}") def clip_video(video_id: str, frm: str = Query("0:00", alias="from"), to: str = Query("", alias="to")): @@ -436,11 +627,31 @@ def _video_dict(v) -> dict: tags = [] return { "video_id": v.video_id, "channel_id": v.channel_id, "title": v.title, "url": v.url, - "upload_date": v.upload_date, "duration": v.duration, "status": v.status, + "upload_date": v.upload_date, + "upload_date_approx": bool(v.upload_date_approx), + "duration": v.duration, "status": v.status, "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, + # 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 + # exact. The sentinel means "newer than anything dated in this channel", + # which is a rank, not a date, so it does not reach the client. + "sort_date": None if v.sort_date == store_mod.NO_DATE_SENTINEL else v.sort_date, + # date_estimated = "la fecha que ves no es exacta": tanto la inferida + # por rank como la aproximada del discovery (texto relativo de YouTube). + "date_estimated": bool( + (v.sort_date and not v.upload_date and v.sort_date != store_mod.NO_DATE_SENTINEL) + or v.upload_date_approx + ), "error_msg": v.error_msg, + "description": v.description, + "availability": v.availability, + # None when fetchable; otherwise members_only / premium_only / private / + # needs_auth / removed. Derived, so it stays correct for rows recorded + # before `availability` existed. + "block_reason": v.block_reason, } diff --git a/src/yt_scraper/webapp/static/app.js b/src/yt_scraper/webapp/static/app.js index 87799df..7e55a75 100644 --- a/src/yt_scraper/webapp/static/app.js +++ b/src/yt_scraper/webapp/static/app.js @@ -999,19 +999,22 @@ }, async deleteJob(id) { + // Recent-jobs table action: always purges the row. Canceling a live + // job is done from the job widget ("Cancel" button there), where the + // request goes through DELETE /api/scrape/{id} without ?action=purge. const job = (this.scrape.jobs || []).find(j => j.id === id); const isTerminal = job && this.jobIsTerminal(job); const msg = isTerminal ? "Remove this finished job from history? Downloaded .md files are not affected." - : "Cancel this running job? Videos already processed stay in your library."; + : "Remove this row from the list? If the job is still running, cancel it from the job widget instead — this only deletes the record."; const ok = await this.confirmDialog({ - title: isTerminal ? "Clear job" : "Cancel job", + title: "Clear job", message: msg, - confirmLabel: isTerminal ? "Remove" : "Cancel job", + confirmLabel: "Remove", danger: !isTerminal, }); if (!ok) return; - try { await this.api("/api/scrape/" + encodeURIComponent(id), { method: "DELETE" }); this.loadJobs(); } + try { await this.api("/api/scrape/" + encodeURIComponent(id) + "?action=purge", { method: "DELETE" }); this.loadJobs(); this.loadDashboard(); } catch (e) { this.toast("Action failed: " + e.message, "error"); } }, @@ -1051,10 +1054,10 @@ const d = await this.api("/api/scrape/history", { method: "DELETE" }); this.toast("Cleared " + ((d && d.deleted) || 0) + " finished job(s)", "success"); this.loadJobs(); + this.loadDashboard(); } catch (e) { this.toast("Clear failed: " + e.message, "error"); } }, jobIsTerminal(j) { return j && ["done", "error", "cancelled"].includes(j.status); }, - jobActionLabel(j) { return this.jobIsTerminal(j) ? "Clear" : "Cancel"; }, scrapePct() { const p = this.scrape.progress || {}; @@ -1261,13 +1264,15 @@ // of a column full of dashes in the middle of a sorted list. videoDate(v) { if (!v) return "—"; - if (v.upload_date) return this.fmtDate(v.upload_date); + if (v.upload_date && !v.upload_date_approx) return this.fmtDate(v.upload_date); + if (v.upload_date) return "~ " + this.fmtDate(v.upload_date); if (v.sort_date) return "~ " + this.fmtDate(v.sort_date); return "—"; }, videoDateTitle(v) { if (!v) return ""; - if (v.upload_date) return "Upload date"; + if (v.upload_date && !v.upload_date_approx) return "Upload date"; + if (v.upload_date_approx) return "Approximate upload date — YouTube shows it as relative text ('3 weeks ago') in listings. Download the .md to learn the exact one."; if (v.date_estimated) return "Not scraped yet — approximate date, taken from the next newer video in the channel. Download the .md to learn the exact one."; return "Not scraped yet — YouTube's channel listing does not report upload dates. Download the .md to learn it."; }, diff --git a/tests/test_approx_dates.py b/tests/test_approx_dates.py new file mode 100644 index 0000000..4bb2dc3 --- /dev/null +++ b/tests/test_approx_dates.py @@ -0,0 +1,93 @@ +"""Fechas aproximadas de discovery: bandera upload_date_approx. + +Matriz de precedencia de upsert_videos + upgrade approx->real de +set_upload_date + conversion timestamp -> YYYYMMDD en discover. +Spec: docs/superpowers/specs/2026-08-22-approximate-upload-dates-design.md +""" + +from yt_scraper.store import Store, VideoRef + + +def _ref(vid, date=None, approx=0, ch="ch1"): + return VideoRef( + video_id=vid, + channel_id=ch, + title="t " + vid, + url=f"https://youtu.be/{vid}", + upload_date=date, + date_approx=approx, + ) + + +def _store(tmp_path): + s = Store(tmp_path / "test.db") + s.upsert_channel("ch1", None, "Canal de prueba") + return s + + +def test_approx_into_empty(tmp_path): + s = _store(tmp_path) + s.upsert_videos([_ref("v1", date="20260820", approx=1)]) + row = s.get_video("v1") + assert row.upload_date == "20260820" + assert row.upload_date_approx == 1 + + +def test_real_not_downgraded_by_approx(tmp_path): + s = _store(tmp_path) + s.upsert_videos([_ref("v1", date="20240101", approx=0)]) + s.upsert_videos([_ref("v1", date="20260820", approx=1)]) + row = s.get_video("v1") + assert row.upload_date == "20240101" + assert row.upload_date_approx == 0 + + +def test_approx_refreshed_by_newer_approx(tmp_path): + s = _store(tmp_path) + s.upsert_videos([_ref("v1", date="20260101", approx=1)]) + s.upsert_videos([_ref("v1", date="20260820", approx=1)]) + row = s.get_video("v1") + assert row.upload_date == "20260820" + assert row.upload_date_approx == 1 + + +def test_extraction_upgrades_approx_to_real(tmp_path): + s = _store(tmp_path) + s.upsert_videos([_ref("v1", date="20260820", approx=1)]) + s.set_upload_date("v1", "20260818") + row = s.get_video("v1") + assert row.upload_date == "20260818" + assert row.upload_date_approx == 0 + + +def test_set_upload_date_still_skips_existing_real(tmp_path): + s = _store(tmp_path) + s.upsert_videos([_ref("v1", date="20240101", approx=0)]) + s.set_upload_date("v1", "20260818") + assert s.get_video("v1").upload_date == "20240101" + + +def test_migration_defaults_flag_to_zero(tmp_path): + # Filas creadas por la via vieja (INSERT sin la columna) quedan con 0. + s = _store(tmp_path) + s.upsert_videos([_ref("v2", date=None, approx=0)]) + row = s.get_video("v2") + assert not row.upload_date_approx + + +# --- conversion en discover ------------------------------------------------- + +from yt_scraper.discover import _entry_upload_date + + +def test_entry_upload_date_prefers_exact(): + assert _entry_upload_date({"upload_date": "20240101"}) == ("20240101", 0) + + +def test_entry_upload_date_from_timestamp(): + ts = 1716115200 # 2024-05-19 12:00 UTC + assert _entry_upload_date({"timestamp": ts}) == ("20240519", 1) + + +def test_entry_upload_date_none(): + assert _entry_upload_date({}) == (None, 0)