feat: fechas aproximadas gratis en discovery (approximate_date) con upgrade a real al extraer
This commit is contained in:
+233
-8
@@ -1,31 +1,76 @@
|
|||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
import logging
|
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
|
import yt_dlp
|
||||||
|
|
||||||
|
from .ratelimit import GLOBAL_PACER, ydl_throttle_opts
|
||||||
from .store import VideoRef
|
from .store import VideoRef
|
||||||
|
|
||||||
log = logging.getLogger(__name__)
|
log = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
|
||||||
def discover_channel(channel_url: str, sleep_subrequests: float = 2.0) -> tuple[str, str, list[VideoRef]]:
|
def _entry_upload_date(entry: dict[str, Any]) -> tuple[str | None, int]:
|
||||||
"""Returns (channel_id, channel_name, video_refs)."""
|
"""(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] = {
|
ydl_opts: dict[str, Any] = {
|
||||||
"extract_flat": "in_playlist",
|
"extract_flat": "in_playlist",
|
||||||
"quiet": True,
|
"quiet": True,
|
||||||
"no_warnings": True,
|
"no_warnings": True,
|
||||||
"skip_download": True,
|
"skip_download": True,
|
||||||
"extract_flat_args": None,
|
# Parse the relative time text ("3 weeks ago") YouTube already includes
|
||||||
"sleep_subrequests": sleep_subrequests,
|
# 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:
|
with yt_dlp.YoutubeDL(ydl_opts) as ydl:
|
||||||
info = ydl.extract_info(channel_url, download=False)
|
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_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"
|
channel_name = info.get("channel") or info.get("title") or info.get("uploader") or "Unknown"
|
||||||
|
avatar = _pick_channel_avatar(info)
|
||||||
|
|
||||||
entries = _flatten_entries(info)
|
entries = _flatten_entries(info)
|
||||||
refs: list[VideoRef] = []
|
refs: list[VideoRef] = []
|
||||||
@@ -36,22 +81,202 @@ def discover_channel(channel_url: str, sleep_subrequests: float = 2.0) -> tuple[
|
|||||||
if entry.get("_type") == "playlist":
|
if entry.get("_type") == "playlist":
|
||||||
continue
|
continue
|
||||||
url = entry.get("url") or f"https://www.youtube.com/watch?v={video_id}"
|
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")
|
duration = entry.get("duration")
|
||||||
title = entry.get("title") or video_id
|
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(
|
refs.append(
|
||||||
VideoRef(
|
VideoRef(
|
||||||
video_id=str(video_id),
|
video_id=str(video_id),
|
||||||
channel_id=str(channel_id),
|
channel_id=str(channel_id),
|
||||||
title=str(title),
|
title=str(title),
|
||||||
url=str(url),
|
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,
|
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)
|
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]]:
|
def _flatten_entries(info: dict[str, Any]) -> list[dict[str, Any]]:
|
||||||
|
|||||||
+528
-50
@@ -49,9 +49,42 @@ _VIDEO_COLUMNS: dict[str, str] = {
|
|||||||
"description": "TEXT",
|
"description": "TEXT",
|
||||||
"chapters_json": "TEXT",
|
"chapters_json": "TEXT",
|
||||||
"segments_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 = """
|
_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 (
|
CREATE TABLE IF NOT EXISTS transcript_segments (
|
||||||
video_id TEXT NOT NULL,
|
video_id TEXT NOT NULL,
|
||||||
idx INTEGER 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
|
@dataclass
|
||||||
class VideoRef:
|
class VideoRef:
|
||||||
video_id: str
|
video_id: str
|
||||||
@@ -100,6 +258,17 @@ class VideoRef:
|
|||||||
url: str
|
url: str
|
||||||
upload_date: str | None = None
|
upload_date: str | None = None
|
||||||
duration: int | 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
|
@dataclass
|
||||||
@@ -123,6 +292,35 @@ class VideoRow:
|
|||||||
description: str | None = None
|
description: str | None = None
|
||||||
chapters_json: str | None = None
|
chapters_json: str | None = None
|
||||||
segments_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
|
@dataclass
|
||||||
@@ -192,7 +390,15 @@ class Store:
|
|||||||
for col, coltype in _VIDEO_COLUMNS.items():
|
for col, coltype in _VIDEO_COLUMNS.items():
|
||||||
if col not in existing:
|
if col not in existing:
|
||||||
conn.execute(f"ALTER TABLE videos ADD COLUMN {col} {coltype}")
|
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)
|
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
|
@contextmanager
|
||||||
def _cursor(self) -> Iterator[sqlite3.Cursor]:
|
def _cursor(self) -> Iterator[sqlite3.Cursor]:
|
||||||
@@ -205,18 +411,19 @@ class Store:
|
|||||||
|
|
||||||
# ------------------------------------------------------------------ channels
|
# ------------------------------------------------------------------ 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()
|
now = _now_iso()
|
||||||
with self._cursor() as cur:
|
with self._cursor() as cur:
|
||||||
cur.execute(
|
cur.execute(
|
||||||
"""INSERT INTO channels (channel_id, handle, name, last_scraped, video_count)
|
"""INSERT INTO channels (channel_id, handle, name, last_scraped, video_count, avatar)
|
||||||
VALUES (?, ?, ?, ?, ?)
|
VALUES (?, ?, ?, ?, ?, ?)
|
||||||
ON CONFLICT(channel_id) DO UPDATE SET
|
ON CONFLICT(channel_id) DO UPDATE SET
|
||||||
handle = excluded.handle,
|
handle = excluded.handle,
|
||||||
name = excluded.name,
|
name = excluded.name,
|
||||||
last_scraped = excluded.last_scraped,
|
last_scraped = excluded.last_scraped,
|
||||||
video_count = excluded.video_count""",
|
video_count = excluded.video_count,
|
||||||
(channel_id, handle, name, now, video_count),
|
avatar = COALESCE(excluded.avatar, channels.avatar)""",
|
||||||
|
(channel_id, handle, name, now, video_count, avatar),
|
||||||
)
|
)
|
||||||
|
|
||||||
def list_channels(self) -> list[dict]:
|
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 videos WHERE channel_id = ?", (channel_id,))
|
||||||
cur.execute("DELETE FROM channels 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
|
# ------------------------------------------------------------------ videos
|
||||||
|
|
||||||
def upsert_videos(self, refs: list[VideoRef]) -> int:
|
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()
|
now = _now_iso()
|
||||||
inserted = 0
|
inserted = 0
|
||||||
with self._cursor() as cur:
|
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:
|
for r in refs:
|
||||||
cur.execute(
|
cur.execute(
|
||||||
"""INSERT INTO videos
|
"""INSERT INTO videos
|
||||||
(video_id, channel_id, title, url, upload_date, duration, status, discovered_at)
|
(video_id, channel_id, title, url, upload_date,
|
||||||
VALUES (?, ?, ?, ?, ?, ?, 'pending', ?)
|
upload_date_approx, duration, availability,
|
||||||
|
channel_seq, status, discovered_at)
|
||||||
|
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, 'pending', ?)
|
||||||
ON CONFLICT(video_id) DO UPDATE SET
|
ON CONFLICT(video_id) DO UPDATE SET
|
||||||
title = excluded.title,
|
title = COALESCE(excluded.title, videos.title),
|
||||||
upload_date = excluded.upload_date,
|
upload_date = CASE
|
||||||
duration = excluded.duration""",
|
WHEN excluded.upload_date IS NULL THEN videos.upload_date
|
||||||
(r.video_id, r.channel_id, r.title, r.url, r.upload_date, r.duration, now),
|
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
|
inserted += 1
|
||||||
|
existing.add(r.video_id)
|
||||||
return inserted
|
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'"
|
sql = "SELECT * FROM videos WHERE status = 'pending'"
|
||||||
params: list[Any] = []
|
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:
|
if channel_id:
|
||||||
sql += " AND channel_id = ?"
|
sql += " AND channel_id = ?"
|
||||||
params.append(channel_id)
|
params.append(channel_id)
|
||||||
@@ -300,6 +643,7 @@ class Store:
|
|||||||
sort: str = "upload_date_desc",
|
sort: str = "upload_date_desc",
|
||||||
page: int = 1,
|
page: int = 1,
|
||||||
size: int = 50,
|
size: int = 50,
|
||||||
|
blocked: bool | None = None,
|
||||||
) -> tuple[list[VideoRow], int]:
|
) -> tuple[list[VideoRow], int]:
|
||||||
where: list[str] = []
|
where: list[str] = []
|
||||||
params: list[Any] = []
|
params: list[Any] = []
|
||||||
@@ -307,6 +651,21 @@ class Store:
|
|||||||
where.append("channel_id = ?"); params.append(channel_id)
|
where.append("channel_id = ?"); params.append(channel_id)
|
||||||
if status:
|
if status:
|
||||||
where.append("status = ?"); params.append(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:
|
if date_from:
|
||||||
where.append("upload_date >= ?"); params.append(date_from.replace("-", ""))
|
where.append("upload_date >= ?"); params.append(date_from.replace("-", ""))
|
||||||
if date_to:
|
if date_to:
|
||||||
@@ -317,19 +676,14 @@ class Store:
|
|||||||
where.append("(title LIKE ? OR description LIKE ?)")
|
where.append("(title LIKE ? OR description LIKE ?)")
|
||||||
params.extend([f"%{q}%", f"%{q}%"])
|
params.extend([f"%{q}%", f"%{q}%"])
|
||||||
clause = ("WHERE " + " AND ".join(where)) if where else ""
|
clause = ("WHERE " + " AND ".join(where)) if where else ""
|
||||||
order = {
|
order = _order_clause(sort)
|
||||||
"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")
|
|
||||||
offset = max(0, (page - 1) * size)
|
offset = max(0, (page - 1) * size)
|
||||||
with self._cursor() as cur:
|
with self._cursor() as cur:
|
||||||
cur.execute(f"SELECT COUNT(*) AS n FROM videos {clause}", params)
|
cur.execute(f"SELECT COUNT(*) AS n FROM videos {clause}", params)
|
||||||
total = cur.fetchone()["n"]
|
total = cur.fetchone()["n"]
|
||||||
cur.execute(
|
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],
|
[*params, size, offset],
|
||||||
)
|
)
|
||||||
rows = [_row_to_videorow(r) for r in cur.fetchall()]
|
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),
|
(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()
|
now = _now_iso()
|
||||||
with self._cursor() as cur:
|
with self._cursor() as cur:
|
||||||
cur.execute(
|
cur.execute(
|
||||||
"UPDATE videos SET status = ?, processed_at = ? WHERE video_id = ?",
|
"UPDATE videos SET status = ?, error_msg = ?, processed_at = ? WHERE video_id = ?",
|
||||||
(status, now, 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:
|
def set_upload_date(self, video_id: str, upload_date: str | None) -> None:
|
||||||
if not upload_date:
|
if not upload_date:
|
||||||
return
|
return
|
||||||
with self._cursor() as cur:
|
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(
|
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),
|
(upload_date, video_id),
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -626,27 +1005,6 @@ class Store:
|
|||||||
r["status"]: r["n"]
|
r["status"]: r["n"]
|
||||||
for r in cur.execute("SELECT status, COUNT(*) AS n FROM videos GROUP BY status").fetchall()
|
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)
|
# top tags (tags is JSON array text)
|
||||||
tag_rows = cur.execute("SELECT tags FROM videos WHERE tags IS NOT NULL AND tags != '[]'").fetchall()
|
tag_rows = cur.execute("SELECT tags FROM videos WHERE tags IS NOT NULL AND tags != '[]'").fetchall()
|
||||||
tag_counts: dict[str, int] = {}
|
tag_counts: dict[str, int] = {}
|
||||||
@@ -662,14 +1020,59 @@ class Store:
|
|||||||
return {
|
return {
|
||||||
"channels": channels,
|
"channels": channels,
|
||||||
"status_breakdown": status_breakdown,
|
"status_breakdown": status_breakdown,
|
||||||
"uploads_over_time": uploads,
|
|
||||||
"duration_histogram": hist,
|
|
||||||
"top_tags": [{"tag": t, "count": c} for t, c in top_tags],
|
"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:
|
def reset_errors(self, channel_id: str | None = None) -> int:
|
||||||
sql = "UPDATE videos SET status = 'pending', error_msg = NULL WHERE status = 'error'"
|
return self.reset_videos(channel_id, statuses=("error",))
|
||||||
params: list[Any] = []
|
|
||||||
|
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:
|
if channel_id:
|
||||||
sql += " AND channel_id = ?"
|
sql += " AND channel_id = ?"
|
||||||
params.append(channel_id)
|
params.append(channel_id)
|
||||||
@@ -677,6 +1080,31 @@ class Store:
|
|||||||
cur.execute(sql, params)
|
cur.execute(sql, params)
|
||||||
return cur.rowcount
|
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:
|
def _sanitize_fts(query: str) -> str:
|
||||||
# Build a safe AND FTS5 query from whitespace-separated terms.
|
# 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")
|
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:
|
def _row_to_videorow(row: sqlite3.Row) -> VideoRow:
|
||||||
keys = row.keys()
|
keys = row.keys()
|
||||||
return VideoRow(
|
return VideoRow(
|
||||||
@@ -717,6 +1191,10 @@ def _row_to_videorow(row: sqlite3.Row) -> VideoRow:
|
|||||||
description=row["description"] if "description" in keys else None,
|
description=row["description"] if "description" in keys else None,
|
||||||
chapters_json=row["chapters_json"] if "chapters_json" 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,
|
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,
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
+221
-10
@@ -12,8 +12,15 @@ from .. import analysis as analysis_mod
|
|||||||
from .. import cookies as cookies_mod
|
from .. import cookies as cookies_mod
|
||||||
from .. import export as export_mod
|
from .. import export as export_mod
|
||||||
from ..config import Config
|
from ..config import Config
|
||||||
|
from ..ratelimit import polite_sleep
|
||||||
|
from .. import store as store_mod
|
||||||
from ..store import Store
|
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:
|
def build_router(store: Store, cfg: Config, jobs) -> APIRouter:
|
||||||
r = APIRouter(prefix="/api")
|
r = APIRouter(prefix="/api")
|
||||||
@@ -28,16 +35,33 @@ def build_router(store: Store, cfg: Config, jobs) -> APIRouter:
|
|||||||
def channels_list():
|
def channels_list():
|
||||||
return {"items": store.list_channels()}
|
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")
|
@r.post("/channels")
|
||||||
async def add_channel(payload: dict):
|
def add_channel(payload: dict):
|
||||||
from ..discover import discover_channel
|
from ..discover import discover_channel, deep_channel_avatar
|
||||||
|
from ..pipeline import cache_channel_avatar
|
||||||
url = (payload or {}).get("url")
|
url = (payload or {}).get("url")
|
||||||
if not url:
|
if not url:
|
||||||
raise HTTPException(400, "url required")
|
raise HTTPException(400, "url required")
|
||||||
cid, name, refs = discover_channel(url, sleep_subrequests=cfg.yt_dlp.sleep_subrequests)
|
cid, name, avatar, refs = discover_channel(url, sleep_subrequests=cfg.yt_dlp.sleep_subrequests)
|
||||||
store.upsert_channel(cid, _handle(url), name, len(refs))
|
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)
|
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}")
|
@r.delete("/channels/{channel_id}")
|
||||||
def del_channel(channel_id: str):
|
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"),
|
date_to: str | None = Query(None, alias="to"),
|
||||||
min_dur: int | None = None, q: str | None = None,
|
min_dur: int | None = None, q: str | None = None,
|
||||||
sort: str = "upload_date_desc", page: int = 1, size: int = 50,
|
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}
|
return {"items": [_video_dict(v) for v in rows], "total": total, "page": page, "size": size}
|
||||||
|
|
||||||
@r.get("/videos/{video_id}")
|
@r.get("/videos/{video_id}")
|
||||||
@@ -140,10 +167,18 @@ def build_router(store: Store, cfg: Config, jobs) -> APIRouter:
|
|||||||
return {"deleted": n}
|
return {"deleted": n}
|
||||||
|
|
||||||
@r.delete("/scrape/{job_id}")
|
@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)
|
job = store.get_job(job_id)
|
||||||
if not job:
|
if not job:
|
||||||
raise HTTPException(404, "job not found")
|
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:
|
if job.status in store.TERMINAL_STATUSES:
|
||||||
store.delete_job(job_id)
|
store.delete_job(job_id)
|
||||||
return {"deleted": job_id}
|
return {"deleted": job_id}
|
||||||
@@ -350,19 +385,31 @@ def build_router(store: Store, cfg: Config, jobs) -> APIRouter:
|
|||||||
|
|
||||||
# -------------------------------------------------- thumbnails (local cache)
|
# -------------------------------------------------- thumbnails (local cache)
|
||||||
@r.post("/tools/thumbnails")
|
@r.post("/tools/thumbnails")
|
||||||
async def download_thumbnails(payload: dict):
|
def download_thumbnails(payload: dict):
|
||||||
from ..pipeline import cache_thumbnail
|
from ..pipeline import cache_thumbnail
|
||||||
video_ids = (payload or {}).get("video_ids") or []
|
video_ids = (payload or {}).get("video_ids") or []
|
||||||
channel_id = (payload or {}).get("channel_id")
|
channel_id = (payload or {}).get("channel_id")
|
||||||
|
explicit = bool(video_ids)
|
||||||
if channel_id and not video_ids:
|
if channel_id and not video_ids:
|
||||||
video_ids = [v.video_id for v in store.get_all(channel_id)]
|
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 = Path(cfg.output_dir_resolved).parent / "thumbnails"
|
||||||
out_dir.mkdir(parents=True, exist_ok=True)
|
out_dir.mkdir(parents=True, exist_ok=True)
|
||||||
n = 0
|
n = 0
|
||||||
for vid in video_ids:
|
for vid in video_ids:
|
||||||
if cache_thumbnail(store, vid, out_dir):
|
if cache_thumbnail(store, vid, out_dir):
|
||||||
n += 1
|
n += 1
|
||||||
return {"downloaded": n, "dir": str(out_dir)}
|
return {"downloaded": n, "dir": str(out_dir), "skipped": skipped}
|
||||||
|
|
||||||
@r.get("/thumbnails/{video_id}")
|
@r.get("/thumbnails/{video_id}")
|
||||||
def serve_thumbnail(video_id: str):
|
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"
|
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)
|
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)
|
# -------------------------------------------------- clip (transcript segment)
|
||||||
@r.get("/clip/{video_id}")
|
@r.get("/clip/{video_id}")
|
||||||
def clip_video(video_id: str, frm: str = Query("0:00", alias="from"), to: str = Query("", alias="to")):
|
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 = []
|
tags = []
|
||||||
return {
|
return {
|
||||||
"video_id": v.video_id, "channel_id": v.channel_id, "title": v.title, "url": v.url,
|
"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),
|
"transcript_lang": v.transcript_lang, "has_chapters": bool(v.has_chapters),
|
||||||
"view_count": v.view_count, "like_count": v.like_count, "tags": tags,
|
"view_count": v.view_count, "like_count": v.like_count, "tags": tags,
|
||||||
"thumbnail": v.thumbnail, "markdown_path": v.markdown_path,
|
"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,
|
"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,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -999,19 +999,22 @@
|
|||||||
},
|
},
|
||||||
|
|
||||||
async deleteJob(id) {
|
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 job = (this.scrape.jobs || []).find(j => j.id === id);
|
||||||
const isTerminal = job && this.jobIsTerminal(job);
|
const isTerminal = job && this.jobIsTerminal(job);
|
||||||
const msg = isTerminal
|
const msg = isTerminal
|
||||||
? "Remove this finished job from history? Downloaded .md files are not affected."
|
? "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({
|
const ok = await this.confirmDialog({
|
||||||
title: isTerminal ? "Clear job" : "Cancel job",
|
title: "Clear job",
|
||||||
message: msg,
|
message: msg,
|
||||||
confirmLabel: isTerminal ? "Remove" : "Cancel job",
|
confirmLabel: "Remove",
|
||||||
danger: !isTerminal,
|
danger: !isTerminal,
|
||||||
});
|
});
|
||||||
if (!ok) return;
|
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"); }
|
catch (e) { this.toast("Action failed: " + e.message, "error"); }
|
||||||
},
|
},
|
||||||
|
|
||||||
@@ -1051,10 +1054,10 @@
|
|||||||
const d = await this.api("/api/scrape/history", { method: "DELETE" });
|
const d = await this.api("/api/scrape/history", { method: "DELETE" });
|
||||||
this.toast("Cleared " + ((d && d.deleted) || 0) + " finished job(s)", "success");
|
this.toast("Cleared " + ((d && d.deleted) || 0) + " finished job(s)", "success");
|
||||||
this.loadJobs();
|
this.loadJobs();
|
||||||
|
this.loadDashboard();
|
||||||
} catch (e) { this.toast("Clear failed: " + e.message, "error"); }
|
} catch (e) { this.toast("Clear failed: " + e.message, "error"); }
|
||||||
},
|
},
|
||||||
jobIsTerminal(j) { return j && ["done", "error", "cancelled"].includes(j.status); },
|
jobIsTerminal(j) { return j && ["done", "error", "cancelled"].includes(j.status); },
|
||||||
jobActionLabel(j) { return this.jobIsTerminal(j) ? "Clear" : "Cancel"; },
|
|
||||||
|
|
||||||
scrapePct() {
|
scrapePct() {
|
||||||
const p = this.scrape.progress || {};
|
const p = this.scrape.progress || {};
|
||||||
@@ -1261,13 +1264,15 @@
|
|||||||
// of a column full of dashes in the middle of a sorted list.
|
// of a column full of dashes in the middle of a sorted list.
|
||||||
videoDate(v) {
|
videoDate(v) {
|
||||||
if (!v) return "—";
|
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);
|
if (v.sort_date) return "~ " + this.fmtDate(v.sort_date);
|
||||||
return "—";
|
return "—";
|
||||||
},
|
},
|
||||||
videoDateTitle(v) {
|
videoDateTitle(v) {
|
||||||
if (!v) return "";
|
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.";
|
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.";
|
return "Not scraped yet — YouTube's channel listing does not report upload dates. Download the .md to learn it.";
|
||||||
},
|
},
|
||||||
|
|||||||
@@ -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)
|
||||||
Reference in New Issue
Block a user