from __future__ import annotations import json import sqlite3 from contextlib import contextmanager from dataclasses import dataclass from datetime import datetime, timezone from pathlib import Path from typing import Any, Iterator SCHEMA = """ CREATE TABLE IF NOT EXISTS channels ( channel_id TEXT PRIMARY KEY, handle TEXT, name TEXT, last_scraped TEXT, video_count INTEGER DEFAULT 0 ); CREATE TABLE IF NOT EXISTS videos ( video_id TEXT PRIMARY KEY, channel_id TEXT NOT NULL, title TEXT, url TEXT NOT NULL, upload_date TEXT, duration INTEGER, status TEXT NOT NULL DEFAULT 'pending', error_msg TEXT, transcript_lang TEXT, transcript_src TEXT, has_chapters INTEGER DEFAULT 0, markdown_path TEXT, discovered_at TEXT NOT NULL, processed_at TEXT, FOREIGN KEY (channel_id) REFERENCES channels(channel_id) ); CREATE INDEX IF NOT EXISTS idx_videos_status ON videos(status); CREATE INDEX IF NOT EXISTS idx_videos_channel ON videos(channel_id); """ # New columns added by the platform migration (idempotent ALTERs) _VIDEO_COLUMNS: dict[str, str] = { "view_count": "INTEGER", "like_count": "INTEGER", "tags": "TEXT", "thumbnail": "TEXT", "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", "video_download_status": "TEXT DEFAULT 'not_downloaded'", "video_path": "TEXT", "video_filename": "TEXT", "video_size": "INTEGER", "video_downloaded_at": "TEXT", "video_error": "TEXT", } _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, start_sec REAL NOT NULL, end_sec REAL NOT NULL, text TEXT NOT NULL, PRIMARY KEY (video_id, idx), FOREIGN KEY (video_id) REFERENCES videos(video_id) ); CREATE INDEX IF NOT EXISTS idx_segments_video ON transcript_segments(video_id); CREATE VIRTUAL TABLE IF NOT EXISTS transcript_fts USING fts5( video_id UNINDEXED, idx UNINDEXED, start_sec UNINDEXED, end_sec UNINDEXED, text ); CREATE TABLE IF NOT EXISTS cookies_meta ( id TEXT PRIMARY KEY, filename TEXT NOT NULL, label TEXT, added_at TEXT NOT NULL, expires_at TEXT, is_active INTEGER DEFAULT 0, has_session INTEGER DEFAULT 0, cookie_count INTEGER DEFAULT 0 ); CREATE TABLE IF NOT EXISTS scrape_jobs ( id TEXT PRIMARY KEY, channel_id TEXT, opts_json TEXT, status TEXT NOT NULL, started_at TEXT, finished_at TEXT, total INTEGER DEFAULT 0, completed INTEGER DEFAULT 0, last_error TEXT ); """ # 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}", # Por nombre de canal (no por UUID): el subquery consulta channels, # que tiene decenas de filas y PK sobre channel_id. "channel": ( "(SELECT c.name FROM channels c WHERE c.channel_id = videos.channel_id) " "COLLATE NOCASE ASC, videos.channel_id, videos.channel_seq DESC, videos.video_id DESC" ), # Agrupa por ciclo de vida: hecho, luego pendientes, luego los que # necesitan atencion (sin subs), errores al final. "status": ( "CASE videos.status WHEN 'done' THEN 0 WHEN 'pending' THEN 1 " "WHEN 'no_subtitles' THEN 2 ELSE 3 END ASC, " + _NEWEST_FIRST ), }.get((sort or "").strip(), _NEWEST_FIRST) @dataclass class VideoRef: video_id: str channel_id: str title: str 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 class VideoRow: video_id: str channel_id: str title: str | None url: str upload_date: str | None duration: int | None status: str error_msg: str | None transcript_lang: str | None transcript_src: str | None has_chapters: int markdown_path: str | None view_count: int | None = None like_count: int | None = None tags: str | None = None thumbnail: str | None = None 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 video_download_status: str | None = None video_path: str | None = None video_filename: str | None = None video_size: int | None = None video_downloaded_at: str | None = None video_error: str | 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 class SegmentRow: video_id: str idx: int start_sec: float end_sec: float text: str @dataclass class SearchHit: video_id: str channel_id: str title: str | None idx: int start_sec: float end_sec: float snippet: str rank: float @dataclass class CookieRow: id: str filename: str label: str | None added_at: str expires_at: str | None is_active: bool has_session: bool cookie_count: int @dataclass class JobRow: id: str channel_id: str | None opts_json: str | None status: str started_at: str | None finished_at: str | None total: int completed: int last_error: str | None class Store: def __init__(self, db_path: str | Path): self.db_path = Path(db_path) self.db_path.parent.mkdir(parents=True, exist_ok=True) self._init_schema() def _connect(self) -> sqlite3.Connection: conn = sqlite3.connect(str(self.db_path)) conn.row_factory = sqlite3.Row conn.execute("PRAGMA journal_mode=WAL") conn.execute("PRAGMA foreign_keys=ON") return conn def _init_schema(self) -> None: with self._connect() as conn: conn.executescript(SCHEMA) # idempotent column adds existing = {row["name"] for row in conn.execute("PRAGMA table_info(videos)")} 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]: conn = self._connect() try: yield conn.cursor() conn.commit() finally: conn.close() # ------------------------------------------------------------------ channels 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, 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, avatar = COALESCE(excluded.avatar, channels.avatar)""", (channel_id, handle, name, now, video_count, avatar), ) def list_channels(self) -> list[dict]: with self._cursor() as cur: cur.execute("SELECT * FROM channels ORDER BY name") return [dict(r) for r in cur.fetchall()] def get_channel(self, channel_id: str) -> dict | None: with self._cursor() as cur: cur.execute("SELECT * FROM channels WHERE channel_id = ?", (channel_id,)) row = cur.fetchone() return dict(row) if row else None def delete_channel(self, channel_id: str) -> None: with self._cursor() as cur: cur.execute("DELETE FROM transcript_segments WHERE video_id IN (SELECT video_id FROM videos WHERE channel_id = ?)", (channel_id,)) cur.execute("DELETE FROM transcript_fts WHERE video_id IN (SELECT video_id 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,)) 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, upload_date_approx, duration, availability, channel_seq, status, discovered_at) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, 'pending', ?) ON CONFLICT(video_id) DO UPDATE SET 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 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, *, 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) sql += " ORDER BY discovered_at ASC" if limit: sql += " LIMIT ?" params.append(limit) with self._cursor() as cur: cur.execute(sql, params) return [_row_to_videorow(row) for row in cur.fetchall()] def get_all(self, channel_id: str | None = None) -> list[VideoRow]: sql = "SELECT * FROM videos" params: list[Any] = [] if channel_id: sql += " WHERE channel_id = ?" params.append(channel_id) sql += " ORDER BY discovered_at ASC" with self._cursor() as cur: cur.execute(sql, params) return [_row_to_videorow(row) for row in cur.fetchall()] def get_video(self, video_id: str) -> VideoRow | None: with self._cursor() as cur: cur.execute("SELECT * FROM videos WHERE video_id = ?", (video_id,)) row = cur.fetchone() return _row_to_videorow(row) if row else None def query_videos( self, channel_id: str | None = None, status: str | None = None, date_from: str | None = None, date_to: str | None = None, min_duration: int | None = None, q: str | None = None, 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] = [] if channel_id: 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: where.append("upload_date <= ?"); params.append(date_to.replace("-", "")) if min_duration is not None: where.append("duration >= ?"); params.append(min_duration) if q: where.append("(title LIKE ? OR description LIKE ?)") params.extend([f"%{q}%", f"%{q}%"]) clause = ("WHERE " + " AND ".join(where)) if where else "" 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 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()] return rows, total def update_video_metadata( self, video_id: str, *, view_count: int | None = None, like_count: int | None = None, tags: list[str] | None = None, thumbnail: str | None = None, description: str | None = None, chapters_json: str | None = None, segments_json: str | None = None, ) -> None: sets: list[str] = [] params: list[Any] = [] if view_count is not None: sets.append("view_count = ?"); params.append(view_count) if like_count is not None: sets.append("like_count = ?"); params.append(like_count) if tags is not None: sets.append("tags = ?"); params.append(json.dumps(tags, ensure_ascii=False)) if thumbnail is not None: sets.append("thumbnail = ?"); params.append(thumbnail) if description is not None: sets.append("description = ?"); params.append(description[:4000]) if chapters_json is not None: sets.append("chapters_json = ?"); params.append(chapters_json) if segments_json is not None: sets.append("segments_json = ?"); params.append(segments_json) if not sets: return params.append(video_id) with self._cursor() as cur: cur.execute(f"UPDATE videos SET {', '.join(sets)} WHERE video_id = ?", params) def mark_done( self, video_id: str, markdown_path: str, transcript_lang: str | None, transcript_src: str | None, has_chapters: bool, ) -> None: now = _now_iso() with self._cursor() as cur: cur.execute( """UPDATE videos SET status = 'done', markdown_path = ?, transcript_lang = ?, transcript_src = ?, has_chapters = ?, processed_at = ?, error_msg = NULL WHERE video_id = ?""", (markdown_path, transcript_lang, transcript_src, int(has_chapters), now, video_id), ) 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 = ?, 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 update_video_download( self, video_id: str, status: str, *, path: str | None = None, filename: str | None = None, size: int | None = None, error: str | None = None, ) -> None: """Persist media-download state independently from transcript state.""" now = _now_iso() if status == "done" else None with self._cursor() as cur: cur.execute( """UPDATE videos SET video_download_status = ?, video_path = COALESCE(?, video_path), video_filename = COALESCE(?, video_filename), video_size = COALESCE(?, video_size), video_downloaded_at = COALESCE(?, video_downloaded_at), video_error = ? WHERE video_id = ?""", (status, path, filename, size, now, error[:2000] if error else None, 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 = ?, upload_date_approx = 0 WHERE video_id = ? AND (upload_date IS NULL OR COALESCE(upload_date_approx, 0) = 1)""", (upload_date, video_id), ) def mark_error(self, video_id: str, error_msg: str) -> None: now = _now_iso() with self._cursor() as cur: cur.execute( """UPDATE videos SET status = 'error', error_msg = ?, processed_at = ? WHERE video_id = ?""", (error_msg[:500], now, video_id), ) # ------------------------------------------------------------------ segments / FTS def store_segments(self, video_id: str, segments: list) -> None: with self._cursor() as cur: cur.execute("DELETE FROM transcript_segments WHERE video_id = ?", (video_id,)) cur.execute("DELETE FROM transcript_fts WHERE video_id = ?", (video_id,)) rows = [] fts_rows = [] for i, seg in enumerate(segments): start = float(getattr(seg, "start", 0.0)) end = float(getattr(seg, "end", start)) text = (getattr(seg, "text", "") or "").strip() if not text: continue rows.append((video_id, i, start, end, text)) fts_rows.append((video_id, i, start, end, text)) if rows: cur.executemany( "INSERT INTO transcript_segments (video_id, idx, start_sec, end_sec, text) VALUES (?, ?, ?, ?, ?)", rows, ) cur.executemany( "INSERT INTO transcript_fts (video_id, idx, start_sec, end_sec, text) VALUES (?, ?, ?, ?, ?)", fts_rows, ) def get_segments(self, video_id: str) -> list[SegmentRow]: with self._cursor() as cur: cur.execute( "SELECT * FROM transcript_segments WHERE video_id = ? ORDER BY idx ASC", (video_id,), ) return [ SegmentRow( video_id=r["video_id"], idx=r["idx"], start_sec=r["start_sec"], end_sec=r["end_sec"], text=r["text"], ) for r in cur.fetchall() ] def has_segments(self, video_id: str) -> bool: with self._cursor() as cur: cur.execute("SELECT 1 FROM transcript_segments WHERE video_id = ? LIMIT 1", (video_id,)) return cur.fetchone() is not None def search_segments(self, query: str, channel_id: str | None = None, limit: int = 50) -> list[SearchHit]: # FTS5 MATCH; join videos for title + optional channel filter. fts_query = _sanitize_fts(query) if not fts_query: return [] sql = ( "SELECT f.video_id, f.idx, f.start_sec, f.end_sec, f.text, v.channel_id, v.title, bm25(transcript_fts) AS rank " "FROM transcript_fts f JOIN videos v ON v.video_id = f.video_id " "WHERE transcript_fts MATCH ?" ) params: list[Any] = [fts_query] if channel_id: sql += " AND v.channel_id = ?" params.append(channel_id) sql += " ORDER BY rank LIMIT ?" params.append(limit) with self._cursor() as cur: cur.execute(sql, params) return [ SearchHit( video_id=r["video_id"], channel_id=r["channel_id"], title=r["title"], idx=r["idx"], start_sec=r["start_sec"], end_sec=r["end_sec"], snippet=r["text"], rank=r["rank"], ) for r in cur.fetchall() ] # ------------------------------------------------------------------ cookies def upsert_cookie(self, cookie_id: str, filename: str, label: str | None, expires_at: str | None, has_session: bool, cookie_count: int) -> None: now = _now_iso() with self._cursor() as cur: cur.execute( """INSERT INTO cookies_meta (id, filename, label, added_at, expires_at, is_active, has_session, cookie_count) VALUES (?, ?, ?, ?, ?, 0, ?, ?) ON CONFLICT(id) DO UPDATE SET filename = excluded.filename, label = excluded.label, expires_at = excluded.expires_at, has_session = excluded.has_session, cookie_count = excluded.cookie_count""", (cookie_id, filename, label, now, expires_at, int(has_session), cookie_count), ) def list_cookies(self) -> list[CookieRow]: with self._cursor() as cur: cur.execute("SELECT * FROM cookies_meta ORDER BY added_at DESC") return [_row_to_cookierow(r) for r in cur.fetchall()] def get_cookie(self, cookie_id: str) -> CookieRow | None: with self._cursor() as cur: cur.execute("SELECT * FROM cookies_meta WHERE id = ?", (cookie_id,)) r = cur.fetchone() return _row_to_cookierow(r) if r else None def set_active_cookie(self, cookie_id: str | None) -> None: with self._cursor() as cur: cur.execute("UPDATE cookies_meta SET is_active = 0") if cookie_id: cur.execute("UPDATE cookies_meta SET is_active = 1 WHERE id = ?", (cookie_id,)) def get_active_cookie(self) -> CookieRow | None: with self._cursor() as cur: cur.execute("SELECT * FROM cookies_meta WHERE is_active = 1 LIMIT 1") r = cur.fetchone() return _row_to_cookierow(r) if r else None def delete_cookie(self, cookie_id: str) -> None: with self._cursor() as cur: cur.execute("DELETE FROM cookies_meta WHERE id = ?", (cookie_id,)) # ------------------------------------------------------------------ jobs def create_job(self, job_id: str, channel_id: str | None, opts: dict) -> None: now = _now_iso() with self._cursor() as cur: cur.execute( """INSERT INTO scrape_jobs (id, channel_id, opts_json, status, started_at, total, completed) VALUES (?, ?, ?, 'queued', ?, 0, 0)""", (job_id, channel_id, json.dumps(opts, ensure_ascii=False), now), ) def update_job(self, job_id: str, *, status: str | None = None, total: int | None = None, completed: int | None = None, last_error: str | None = None, finished: bool = False) -> None: sets: list[str] = [] params: list[Any] = [] if status: sets.append("status = ?"); params.append(status) if total is not None: sets.append("total = ?"); params.append(total) if completed is not None: sets.append("completed = ?"); params.append(completed) if last_error is not None: sets.append("last_error = ?"); params.append(last_error[:500]) if finished: sets.append("finished_at = ?"); params.append(_now_iso()) if not sets: return params.append(job_id) with self._cursor() as cur: cur.execute(f"UPDATE scrape_jobs SET {', '.join(sets)} WHERE id = ?", params) def get_job(self, job_id: str) -> JobRow | None: with self._cursor() as cur: cur.execute("SELECT * FROM scrape_jobs WHERE id = ?", (job_id,)) r = cur.fetchone() return _row_to_jobrow(r) if r else None def list_jobs(self, limit: int = 50) -> list[JobRow]: with self._cursor() as cur: cur.execute("SELECT * FROM scrape_jobs ORDER BY started_at DESC LIMIT ?", (limit,)) return [_row_to_jobrow(r) for r in cur.fetchall()] TERMINAL_STATUSES = ("done", "error", "cancelled") def delete_job(self, job_id: str) -> bool: with self._cursor() as cur: cur.execute("DELETE FROM scrape_jobs WHERE id = ?", (job_id,)) return cur.rowcount > 0 def delete_terminal_jobs(self) -> int: placeholders = ",".join("?" for _ in self.TERMINAL_STATUSES) with self._cursor() as cur: cur.execute( f"DELETE FROM scrape_jobs WHERE status IN ({placeholders})", list(self.TERMINAL_STATUSES), ) return cur.rowcount # ------------------------------------------------------------------ aggregates def stats(self, channel_id: str | None = None) -> dict[str, int]: sql = "SELECT status, COUNT(*) as n FROM videos" params: list[Any] = [] if channel_id: sql += " WHERE channel_id = ?" params.append(channel_id) sql += " GROUP BY status" result: dict[str, int] = {} with self._cursor() as cur: cur.execute(sql, params) for row in cur.fetchall(): result[row["status"]] = row["n"] return result def dashboard(self) -> dict: with self._cursor() as cur: channels = [dict(r) for r in cur.execute("SELECT * FROM channels ORDER BY name").fetchall()] for ch in channels: cid = ch["channel_id"] row = cur.execute( "SELECT COUNT(*) AS n, COALESCE(SUM(duration),0) AS dur, MIN(upload_date) AS mind, MAX(upload_date) AS maxd, COALESCE(SUM(view_count),0) AS views, COALESCE(SUM(like_count),0) AS likes FROM videos WHERE channel_id = ?", (cid,), ).fetchone() ch["video_count_db"] = row["n"] ch["total_duration"] = row["dur"] ch["date_min"] = row["mind"] ch["date_max"] = row["maxd"] ch["total_views"] = row["views"] ch["total_likes"] = row["likes"] status_breakdown = { r["status"]: r["n"] for r in cur.execute("SELECT status, COUNT(*) AS n FROM videos GROUP BY status").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] = {} for tr in tag_rows: try: for t in json.loads(tr["tags"]): t = (t or "").strip().lower() if t: tag_counts[t] = tag_counts.get(t, 0) + 1 except (json.JSONDecodeError, TypeError): continue top_tags = sorted(tag_counts.items(), key=lambda x: x[1], reverse=True)[:20] return { "channels": channels, "status_breakdown": status_breakdown, "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: 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) with self._cursor() as cur: 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. cleaned = (query or "").strip() if not cleaned: return "" terms = [] for tok in cleaned.split(): tok = tok.strip('"') if tok: terms.append(f'"{tok.replace("\"", "")}"') return " AND ".join(terms) 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( video_id=row["video_id"], channel_id=row["channel_id"], title=row["title"], url=row["url"], upload_date=row["upload_date"], duration=row["duration"], status=row["status"], error_msg=row["error_msg"], transcript_lang=row["transcript_lang"], transcript_src=row["transcript_src"], has_chapters=row["has_chapters"], markdown_path=row["markdown_path"], view_count=row["view_count"] if "view_count" in keys else None, like_count=row["like_count"] if "like_count" in keys else None, tags=row["tags"] if "tags" in keys else None, thumbnail=row["thumbnail"] if "thumbnail" 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, 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, video_download_status=row["video_download_status"] if "video_download_status" in keys else None, video_path=row["video_path"] if "video_path" in keys else None, video_filename=row["video_filename"] if "video_filename" in keys else None, video_size=row["video_size"] if "video_size" in keys else None, video_downloaded_at=row["video_downloaded_at"] if "video_downloaded_at" in keys else None, video_error=row["video_error"] if "video_error" in keys else None, sort_date=row["sort_date"] if "sort_date" in keys else None, ) def _row_to_cookierow(row: sqlite3.Row) -> CookieRow: return CookieRow( id=row["id"], filename=row["filename"], label=row["label"], added_at=row["added_at"], expires_at=row["expires_at"], is_active=bool(row["is_active"]), has_session=bool(row["has_session"]), cookie_count=row["cookie_count"], ) def _row_to_jobrow(row: sqlite3.Row) -> JobRow: return JobRow( id=row["id"], channel_id=row["channel_id"], opts_json=row["opts_json"], status=row["status"], started_at=row["started_at"], finished_at=row["finished_at"], total=row["total"], completed=row["completed"], last_error=row["last_error"], )