#!/usr/bin/env python3 """YouTube Suggester engine for the Omarchy Quattro plugin. Pipeline: 1. Build the watched set from local browser history (Firefox/Chrome family) plus videos previously opened through this plugin. 2. Fetch the logged-in YouTube subscriptions feed via yt-dlp with browser cookies. 3. Filter out watched videos, keep the newest N candidates. 4. Score video metadata against the user's interest keywords and pick the top recommendations with a description snippet. Every state change is written atomically to the state file and printed as a single JSON line, so the QML service can stream live progress. """ import argparse from contextlib import contextmanager import ctypes import hashlib import hmac import json import os import re import selectors import secrets import signal import stat import sqlite3 import subprocess import sys import tempfile import threading import time import webbrowser from concurrent.futures import ThreadPoolExecutor, as_completed from pathlib import Path from urllib.parse import parse_qs, urlsplit HOME = Path.home() CONFIG_DIR = HOME / ".config" / "youtube-suggester" CACHE_DIR = HOME / ".cache" / "omarchy" / "youtube-suggester" CONFIG_FILE = CONFIG_DIR / "config.json" STATE_FILE = CACHE_DIR / "state.json" SEEN_FILE = CACHE_DIR / "seen.json" DISMISSED_FILE = CACHE_DIR / "dismissed.json" DEFAULT_CONFIG = { "interests": ["AI", "Software"], "browser": "chromium", "feed_limit": 120, "use_account_history": True, "account_history_limit": 200, "max_candidates": 60, "metadata_workers": 4, } WATCH_PATTERNS = ("%youtube.com/watch%", "%youtu.be/%", "%youtube.com/shorts%") # Chromium-family cookie decryption support. # Chromium >= ~137 writes v11 cookies on Linux whose plaintext carries a # 32-byte SHA-256 hash prefix; the AES key is derived (PBKDF2-SHA1, # salt "saltysalt", 1 iteration, 16 bytes) from the browser's "Safe Storage" # secret held in the desktop keyring. yt-dlp does not strip that prefix on # Linux yet, so we decrypt the cookie database ourselves and hand yt-dlp a # Netscape cookies.txt file instead. try: from Cryptodome.Cipher import AES as _AES except ImportError: # pragma: no cover try: from Crypto.Cipher import AES as _AES except ImportError: _AES = None CHROMIUM_BROWSERS = { "chromium": { "roots": [HOME / ".config" / "chromium"], "label": "Chromium", "app": "chromium", }, "chrome": { "roots": [HOME / ".config" / "google-chrome"], "label": "Chrome", "app": "Chrome", }, "brave": { "roots": [HOME / ".config" / "BraveSoftware" / "Brave-Browser"], "label": "Brave", "app": "brave", }, "edge": { "roots": [HOME / ".config" / "microsoft-edge"], "label": "Microsoft Edge", "app": "Microsoft Edge", }, } LEGACY_COOKIES_FILE = CACHE_DIR / "cookies.txt" MAX_PERSISTED_JSON_BYTES = 1024 * 1024 MAX_STATE_BYTES = 1024 * 1024 MAX_SUBPROCESS_STDOUT_BYTES = 4 * 1024 * 1024 MAX_SUBPROCESS_STDERR_BYTES = 64 * 1024 MAX_TITLE_CHARS = 500 MAX_CHANNEL_CHARS = 200 MAX_DESCRIPTION_CHARS = 4000 MAX_TAG_CHARS = 100 MAX_TAGS = 12 MAX_STATE_RECORDS = 300 MAX_SEEN_IDS = 5000 MAX_DETAIL_CHARS = 500 MAX_ERROR_CHARS = 2000 MAX_COOKIE_DB_BYTES = 64 * 1024 * 1024 MAX_HISTORY_DB_BYTES = 256 * 1024 * 1024 MAX_WAL_BYTES = 128 * 1024 * 1024 MAX_SHM_BYTES = 16 * 1024 * 1024 MAX_BROWSER_PROFILES = 32 MAX_BROWSER_DIRECTORY_ENTRIES = 128 MAX_COOKIE_ROWS = 10000 MAX_HISTORY_ROWS = 100000 MAX_HISTORY_IDS = 100000 MAX_SQLITE_VALUE_BYTES = 1024 * 1024 MAX_COOKIE_FIELD_BYTES = 64 * 1024 MAX_COOKIE_PAYLOAD_BYTES = 4 * 1024 * 1024 SQLITE_QUERY_TIMEOUT = 10 VIDEO_ID_PATTERN = re.compile(r"[A-Za-z0-9_-]{11}") COOKIE_DOMAINS = ("google.com", "youtube.com") def _clip_text(value, max_chars, default=""): if value is None: return default if not isinstance(value, str): value = str(value) return value[:max_chars] def _bounded_int(value, default, minimum, maximum): try: if isinstance(value, str): value = value[:32] return max(minimum, min(maximum, int(value))) except (OverflowError, TypeError, ValueError): return default def _supervise_with_parent_death(expected_parent, argv): """Keep a child process group tied to this engine's lifetime.""" if not argv: raise SystemExit("missing child command") child = None def stop_child_group(): if child is None: return try: os.killpg(child.pid, signal.SIGKILL) except ProcessLookupError: pass def parent_gone(received_signal, _frame=None): stop_child_group() os._exit(128 + received_signal) signal.signal(signal.SIGTERM, parent_gone) signal.signal(signal.SIGINT, parent_gone) try: libc = ctypes.CDLL(None, use_errno=True) if libc.prctl(1, signal.SIGTERM, 0, 0, 0) != 0: # PR_SET_PDEATHSIG raise OSError(ctypes.get_errno(), "prctl(PR_SET_PDEATHSIG) failed") except (AttributeError, OSError) as exc: raise SystemExit(str(exc)) if os.getppid() != expected_parent: parent_gone(signal.SIGTERM) child = subprocess.Popen(argv, close_fds=False, start_new_session=True) try: returncode = child.wait() finally: stop_child_group() raise SystemExit(returncode) def run_bounded(cmd, timeout, stdout_limit=MAX_SUBPROCESS_STDOUT_BYTES, stderr_limit=MAX_SUBPROCESS_STDERR_BYTES, pass_fds=()): """Run a child while draining both pipes into fixed-size buffers.""" wrapped = [ sys.executable, os.path.realpath(__file__), "_exec-child", str(os.getpid()), *cmd, ] process = subprocess.Popen( wrapped, stdout=subprocess.PIPE, stderr=subprocess.PIPE, start_new_session=True, pass_fds=tuple(pass_fds), ) selector = None buffers = {"stdout": bytearray(), "stderr": bytearray()} limits = {"stdout": stdout_limit, "stderr": stderr_limit} limited = False stopped = False deadline = time.monotonic() + timeout def stop_process(): nonlocal stopped if stopped: return stopped = True try: os.kill(process.pid, signal.SIGTERM) except ProcessLookupError: pass except OSError: try: process.terminate() except ProcessLookupError: pass returncode = None try: selector = selectors.DefaultSelector() selector.register(process.stdout, selectors.EVENT_READ, "stdout") selector.register(process.stderr, selectors.EVENT_READ, "stderr") while selector.get_map(): if stopped: wait = 0.1 else: remaining = deadline - time.monotonic() if remaining <= 0: limited = True stop_process() continue wait = remaining events = selector.select(wait) if not events: if not stopped: limited = True stop_process() continue for key, _ in events: try: chunk = os.read(key.fd, 65536) except OSError: selector.unregister(key.fileobj) continue if not chunk: selector.unregister(key.fileobj) continue stream = key.data remaining = limits[stream] - len(buffers[stream]) if len(chunk) > remaining: if remaining > 0: buffers[stream].extend(chunk[:remaining]) limited = True stop_process() else: buffers[stream].extend(chunk) returncode = process.wait() finally: if process.poll() is None: stop_process() try: process.wait(timeout=5) except subprocess.TimeoutExpired: process.kill() process.wait() if selector is not None: selector.close() for stream in (process.stdout, process.stderr): if stream is not None: stream.close() result = subprocess.CompletedProcess( cmd, returncode, stdout=bytes(buffers["stdout"]).decode("utf-8", errors="replace"), stderr=bytes(buffers["stderr"]).decode("utf-8", errors="replace"), ) return result, limited def _bounded_string_list(value, limit, max_chars): if not isinstance(value, (list, tuple)): return [] return [ _clip_text(item, max_chars) for item in value[:limit] if item is not None ] def _bounded_tags(value): return [tag for tag in _bounded_string_list(value, MAX_TAGS, MAX_TAG_CHARS) if tag] def _bounded_item(item): """Keep only the bounded fields the QML UI can render.""" if not isinstance(item, dict): return None scores = item.get("keyword_scores") bounded_scores = {} if isinstance(scores, dict): for index, (key, value) in enumerate(scores.items()): if index >= MAX_TAGS: break key = _clip_text(key, MAX_TAG_CHARS) if key: bounded_scores[key] = _bounded_int(value, 0, 0, 10**9) return { "id": _clip_text(item.get("id"), 32), "title": _clip_text(item.get("title"), MAX_TITLE_CHARS, "(untitled)"), "channel": _clip_text(item.get("channel"), MAX_CHANNEL_CHARS, "YouTube"), "duration": _bounded_int(item.get("duration"), 0, 0, 10**7), "duration_formatted": _clip_text(item.get("duration_formatted"), 32, "?"), "views": _bounded_int(item.get("views"), 0, 0, 10**18), "score": _bounded_int(item.get("score"), 0, 0, 10**9), "matched": _bounded_string_list(item.get("matched"), MAX_TAGS, MAX_TAG_CHARS), "keyword_scores": bounded_scores, "meta_description": _clip_text( item.get("meta_description"), MAX_DESCRIPTION_CHARS ), "tags": _bounded_tags(item.get("tags")), "thumbnail": _clip_text(item.get("thumbnail"), 500), "url": _clip_text(item.get("url"), 500), "description": _clip_text(item.get("description"), MAX_DESCRIPTION_CHARS), "tag": _clip_text(item.get("tag"), MAX_TAG_CHARS), "source": _clip_text(item.get("source"), 16), "uploaded_at": _bounded_int(item.get("uploaded_at"), 0, 0, 10**12), "age_label": _clip_text(item.get("age_label"), 32), "views_per_hour": _bounded_int(item.get("views_per_hour"), 0, 0, 10**15), "opened_at": _clip_text(item.get("opened_at"), 64), } def _bounded_items(value, limit): if not isinstance(value, list): return [] items = [] for item in value[:limit]: bounded = _bounded_item(item) if bounded is not None: items.append(bounded) return items def _pbkdf2_key(password, length=16): return hashlib.pbkdf2_hmac("sha1", password, b"saltysalt", 1, dklen=length) def _keyring_safe_storage_secret(browser): """Fetch the browser's Safe Storage secret from the desktop keyring.""" info = CHROMIUM_BROWSERS.get(browser) if not info: return None try: import secretstorage con = secretstorage.dbus_init() col = secretstorage.get_default_collection(con) want_label = info["label"] + " Safe Storage" want_app = info["app"].lower() for item in col.get_all_items(): if item.get_label() != want_label: continue attrs = item.get_attributes() if str(attrs.get("application", "")).lower() == want_app: return item.get_secret() except Exception: pass return None def _open_directory_path(path): """Open an absolute directory without following any path component.""" path = Path(path) if not path.is_absolute(): raise ValueError("browser directory must be absolute") flags = os.O_RDONLY | os.O_NOFOLLOW flags |= getattr(os, "O_DIRECTORY", 0) | getattr(os, "O_CLOEXEC", 0) fd = os.open("/", flags) try: for part in path.parts[1:]: next_fd = os.open(part, flags, dir_fd=fd) os.close(fd) fd = next_fd info = os.fstat(fd) if not stat.S_ISDIR(info.st_mode) or info.st_uid != os.getuid(): raise PermissionError("browser directory must be user-owned") return fd except Exception: os.close(fd) raise def _open_profile_directory(root, profile): root_fd = _open_directory_path(root) flags = os.O_RDONLY | os.O_NOFOLLOW flags |= getattr(os, "O_DIRECTORY", 0) | getattr(os, "O_CLOEXEC", 0) try: profile_fd = os.open(profile, flags, dir_fd=root_fd) finally: os.close(root_fd) try: info = os.fstat(profile_fd) if not stat.S_ISDIR(info.st_mode) or info.st_uid != os.getuid(): raise PermissionError("browser profile must be user-owned") except Exception: os.close(profile_fd) raise return profile_fd def _open_bounded_regular_at(directory_fd, name, max_bytes, required=False): flags = os.O_RDONLY | os.O_NONBLOCK | os.O_NOFOLLOW flags |= getattr(os, "O_CLOEXEC", 0) try: fd = os.open(name, flags, dir_fd=directory_fd) except FileNotFoundError: if required: raise return None try: info = os.fstat(fd) if not stat.S_ISREG(info.st_mode) or info.st_uid != os.getuid(): raise PermissionError(f"unsafe browser file: {name}") if info.st_size < 0 or info.st_size > max_bytes: raise ValueError(f"browser file exceeds limit: {name}") return fd except Exception: os.close(fd) raise def _profile_db_candidates(root, db_name, max_bytes): """Return bounded direct-child browser DB candidates under one root.""" try: root_fd = _open_directory_path(root) except OSError: return [] candidates = [] scanned = 0 try: with os.scandir(root_fd) as entries: for entry in entries: scanned += 1 if ( scanned > MAX_BROWSER_DIRECTORY_ENTRIES or len(candidates) >= MAX_BROWSER_PROFILES ): break try: if not entry.is_dir(follow_symlinks=False): continue profile_fd = os.open( entry.name, os.O_RDONLY | os.O_NOFOLLOW | getattr(os, "O_DIRECTORY", 0) | getattr(os, "O_CLOEXEC", 0), dir_fd=root_fd, ) try: db_fd = _open_bounded_regular_at( profile_fd, db_name, max_bytes, required=True ) try: mtime = os.fstat(db_fd).st_mtime finally: os.close(db_fd) finally: os.close(profile_fd) candidates.append((Path(root), entry.name, db_name, mtime)) except (OSError, ValueError): continue finally: os.close(root_fd) return candidates def _find_browser_cookie_db(browser): info = CHROMIUM_BROWSERS.get(browser) if not info: return None candidates = [] for root in info["roots"]: remaining = MAX_BROWSER_PROFILES - len(candidates) candidates.extend( _profile_db_candidates(root, "Cookies", MAX_COOKIE_DB_BYTES)[:remaining] ) if len(candidates) >= MAX_BROWSER_PROFILES: break for candidate in candidates: if candidate[1] == "Default": return candidate return max(candidates, key=lambda item: item[3]) if candidates else None def _copy_fd_bounded(source_fd, destination, max_bytes): destination_fd = os.open( destination, os.O_WRONLY | os.O_CREAT | os.O_EXCL | os.O_NOFOLLOW | getattr(os, "O_CLOEXEC", 0), 0o600, ) copied = 0 try: while True: chunk = os.read(source_fd, min(1024 * 1024, max_bytes - copied + 1)) if not chunk: break copied += len(chunk) if copied > max_bytes: raise ValueError("browser file grew beyond its size limit") view = memoryview(chunk) while view: written = os.write(destination_fd, view) if written <= 0: raise OSError("failed to copy browser file") view = view[written:] except Exception: os.close(destination_fd) try: os.unlink(destination) except OSError: pass raise os.close(destination_fd) def _file_identity(info): return (info.st_dev, info.st_ino, info.st_size, info.st_mtime_ns) def _runtime_temp_parent(): value = os.environ.get("XDG_RUNTIME_DIR") if not value: return None try: fd = _open_directory_path(value) except (OSError, ValueError): return None os.close(fd) return value def _recover_private_snapshot(path, max_bytes): """Apply copied WAL/journal data before reopening the snapshot read-only.""" con = sqlite3.connect(path, timeout=2) deadline = time.monotonic() + SQLITE_QUERY_TIMEOUT try: con.execute("PRAGMA trusted_schema=OFF") if hasattr(con, "setlimit") and hasattr(sqlite3, "SQLITE_LIMIT_LENGTH"): con.setlimit(sqlite3.SQLITE_LIMIT_LENGTH, MAX_SQLITE_VALUE_BYTES) con.set_progress_handler( lambda: 1 if time.monotonic() > deadline else 0, 10000 ) con.execute("PRAGMA schema_version").fetchone() finally: con.close() if path.stat().st_size > max_bytes: raise RuntimeError("recovered browser database exceeds its limit") @contextmanager def _browser_db_snapshot(source, main_limit): """Create a bounded snapshot, retrying if SQLite files change mid-copy.""" root, profile, db_name, _mtime = source with tempfile.TemporaryDirectory( prefix="omarchy-youtube-suggester-", dir=_runtime_temp_parent() ) as temp_dir: profile_fd = _open_profile_directory(root, profile) try: destination = Path(temp_dir) / db_name file_limits = ( (db_name, main_limit, True, True), (db_name + "-wal", MAX_WAL_BYTES, False, True), (db_name + "-journal", MAX_WAL_BYTES, False, True), (db_name + "-shm", MAX_SHM_BYTES, False, False), ) for _attempt in range(3): records = [] try: for name, limit, required, should_copy in file_limits: fd = _open_bounded_regular_at( profile_fd, name, limit, required=required ) identity = _file_identity(os.fstat(fd)) if fd is not None else None records.append((name, limit, should_copy, fd, identity)) for name, limit, should_copy, fd, _identity in records: if fd is not None and should_copy: target = destination if name == db_name else Path(temp_dir) / name _copy_fd_bounded(fd, target, limit) stable = True for name, limit, _should_copy, fd, identity in records: if fd is not None and _file_identity(os.fstat(fd)) != identity: stable = False current_fd = _open_bounded_regular_at(profile_fd, name, limit) try: current = ( _file_identity(os.fstat(current_fd)) if current_fd is not None else None ) finally: if current_fd is not None: os.close(current_fd) if current != identity: stable = False if stable: _recover_private_snapshot(destination, main_limit) yield destination return finally: for _name, _limit, _copy, fd, _identity in records: if fd is not None: os.close(fd) for path in Path(temp_dir).iterdir(): path.unlink() raise RuntimeError("browser database changed during snapshot") finally: os.close(profile_fd) @contextmanager def _query_browser_snapshot(source, main_limit): with _browser_db_snapshot(source, main_limit) as snapshot: con = sqlite3.connect(snapshot.as_uri() + "?mode=ro", uri=True) deadline = time.monotonic() + SQLITE_QUERY_TIMEOUT try: con.execute("PRAGMA query_only=ON") con.execute("PRAGMA trusted_schema=OFF") if hasattr(con, "setlimit") and hasattr(sqlite3, "SQLITE_LIMIT_LENGTH"): con.setlimit(sqlite3.SQLITE_LIMIT_LENGTH, MAX_SQLITE_VALUE_BYTES) con.set_progress_handler( lambda: 1 if time.monotonic() > deadline else 0, 10000 ) yield con finally: con.close() def _decrypt_chromium_cookie(blob, os_secret, host_key, has_host_hash): if not blob: return None version, body = blob[:3], blob[3:] if version == b"v10": attempts = [b"peanuts"] elif version == b"v11": attempts = [] if os_secret: attempts.append(os_secret) attempts.append(b"peanuts") else: return None for password in attempts: try: dec = _AES.new(_pbkdf2_key(password), _AES.MODE_CBC, b" " * 16).decrypt( body ) pad = dec[-1] if not (isinstance(pad, int) and 1 <= pad <= 16): continue if dec[-pad:] != bytes([pad]) * pad: continue dec = dec[:-pad] if has_host_hash: if len(dec) <= 32: continue expected = hashlib.sha256(host_key.encode("utf-8")).digest() if not hmac.compare_digest(dec[:32], expected): continue dec = dec[32:] return dec.decode("utf-8") except Exception: continue return None def _webkit_to_unix(microseconds): if not microseconds: return 0 return int(microseconds / 1_000_000 - 11644473600) def _safe_cookie_field(value): return ( isinstance(value, str) and len(value.encode("utf-8")) <= MAX_COOKIE_FIELD_BYTES and all(ord(char) >= 0x20 and ord(char) != 0x7F for char in value) ) def _cookie_row_line(row, secret, has_host_hash=False): try: host, name, value, encrypted, path, expires, secure = row if not isinstance(host, str): return None normalized = host.lstrip(".").lower() if not any( normalized == domain or normalized.endswith("." + domain) for domain in COOKIE_DOMAINS ): return None if not value: value = _decrypt_chromium_cookie( encrypted, secret, host, has_host_hash ) path = path or "/" if not all(_safe_cookie_field(field) for field in (host, name, value, path)): return None return "\t".join( [ host, "TRUE" if host.startswith(".") else "FALSE", path, "TRUE" if secure else "FALSE", str(max(0, _webkit_to_unix(expires))), name, value, ] ) except (OverflowError, TypeError, ValueError): return None def export_browser_cookies(browser): """Return a bounded Netscape cookie payload without persisting plaintext.""" if _AES is None or browser not in CHROMIUM_BROWSERS: return None source = _find_browser_cookie_db(browser) if not source: return None secret = _keyring_safe_storage_secret(browser) lines = ["# Netscape HTTP Cookie File", "# Generated by omarchy-youtube-suggester"] count = 0 payload_bytes = sum(len(line.encode("utf-8")) + 1 for line in lines) try: with _query_browser_snapshot(source, MAX_COOKIE_DB_BYTES) as con: try: meta_row = con.execute( "SELECT value FROM meta WHERE key='version' LIMIT 1" ).fetchone() except sqlite3.Error: meta_row = None meta_version = _bounded_int( meta_row[0] if meta_row else 0, 0, 0, 1000000 ) has_host_hash = meta_version >= 24 rows = con.execute( "SELECT host_key, name, value, encrypted_value, path, " "expires_utc, is_secure FROM cookies " "WHERE host_key LIKE '%youtube%' OR host_key LIKE '%google%' " "LIMIT ?", (MAX_COOKIE_ROWS + 1,), ) for index, row in enumerate(rows): if index >= MAX_COOKIE_ROWS: return None line = _cookie_row_line(row, secret, has_host_hash) if line is None: continue line_bytes = len(line.encode("utf-8")) + 1 if payload_bytes + line_bytes > MAX_COOKIE_PAYLOAD_BYTES: return None payload_bytes += line_bytes lines.append(line) count += 1 except (OSError, RuntimeError, sqlite3.Error, TypeError, ValueError): return None if count < 5: # a logged-in session always carries many google/youtube cookies return None return ("\n".join(lines) + "\n").encode("utf-8") # --------------------------------------------------------------------------- # persistence helpers _PRIVATE_DIR_FDS = {} def _private_directory_fd(directory): """Open and validate a private persistence directory once per process.""" directory = Path(directory) key = os.path.abspath(os.fspath(directory)) cached = _PRIVATE_DIR_FDS.get(key) if cached is not None: return cached flags = os.O_RDONLY | os.O_NOFOLLOW flags |= getattr(os, "O_DIRECTORY", 0) | getattr(os, "O_CLOEXEC", 0) fd = os.open(os.fspath(directory), flags) try: info = os.fstat(fd) if not stat.S_ISDIR(info.st_mode) or info.st_uid != os.getuid(): raise PermissionError("persistence directory must be a user-owned directory") os.fchmod(fd, 0o700) except Exception: os.close(fd) raise _PRIVATE_DIR_FDS[key] = fd return fd def _private_file_parts(path): path = Path(path) return _private_directory_fd(path.parent), path.name def _open_private_temp(directory_fd, prefix, suffix): """Create a random mode-0600 file beneath a verified directory fd.""" flags = os.O_RDWR | os.O_CREAT | os.O_EXCL | os.O_NOFOLLOW flags |= getattr(os, "O_CLOEXEC", 0) for _ in range(16): name = f"{prefix}{secrets.token_hex(16)}{suffix}" try: return os.open(name, flags, 0o600, dir_fd=directory_fd), name except FileExistsError: continue raise FileExistsError("could not allocate a private temporary file") def ensure_dirs(): for directory in (CONFIG_DIR, CACHE_DIR): directory.mkdir(parents=True, exist_ok=True) _private_directory_fd(directory) cache_fd = _private_directory_fd(CACHE_DIR) try: os.unlink(LEGACY_COOKIES_FILE.name, dir_fd=cache_fd) except FileNotFoundError: pass with os.scandir(cache_fd) as entries: for entry in entries: if not (entry.name.startswith(".cookies-") and entry.name.endswith(".txt")): continue try: info = entry.stat(follow_symlinks=False) if stat.S_ISREG(info.st_mode) and info.st_uid == os.getuid(): os.unlink(entry.name, dir_fd=cache_fd) except FileNotFoundError: pass def _read_bounded_regular_file(path, max_bytes): """Read a caller-owned regular file through a bounded, fixed descriptor.""" limit = _bounded_int( max_bytes, MAX_PERSISTED_JSON_BYTES, 0, MAX_PERSISTED_JSON_BYTES ) flags = os.O_RDONLY | os.O_NONBLOCK | os.O_NOFOLLOW flags |= getattr(os, "O_CLOEXEC", 0) try: directory_fd, name = _private_file_parts(path) fd = os.open(name, flags, dir_fd=directory_fd) except OSError: return None try: info = os.fstat(fd) if not stat.S_ISREG(info.st_mode) or info.st_uid != os.getuid(): return None raw = os.read(fd, limit + 1) if len(raw) > limit: return None return raw except OSError: return None finally: os.close(fd) def load_json(path, default, max_bytes=MAX_PERSISTED_JSON_BYTES): """Read at most one bounded, caller-owned regular JSON file.""" raw = _read_bounded_regular_file(path, max_bytes) if raw is None: return default try: return json.loads(raw.decode("utf-8")) except (UnicodeDecodeError, ValueError, TypeError): return default def _json_bytes(data, indent=None): separators = None if indent is not None else (",", ":") return json.dumps( data, ensure_ascii=True, indent=indent, separators=separators, ).encode("utf-8") def emit_json(data, max_bytes=MAX_PERSISTED_JSON_BYTES): payload = _json_bytes(data) if len(payload) > max_bytes: raise ValueError("JSON output exceeds the configured size limit") print(payload.decode("utf-8"), flush=True) def atomic_write_json(path, data, max_bytes=MAX_PERSISTED_JSON_BYTES): payload = _json_bytes(data, indent=2) if len(payload) > max_bytes: raise ValueError("JSON state exceeds the configured size limit") path = Path(path) directory_fd, name = _private_file_parts(path) tmp_name = None open_fd = None try: tmp_fd, tmp_name = _open_private_temp(directory_fd, f".{name}-", ".tmp") open_fd = tmp_fd os.fchmod(tmp_fd, 0o600) with os.fdopen(tmp_fd, "wb") as fh: open_fd = None fh.write(payload) fh.flush() os.fsync(fh.fileno()) os.replace( tmp_name, name, src_dir_fd=directory_fd, dst_dir_fd=directory_fd, ) tmp_name = None finally: if open_fd is not None: try: os.close(open_fd) except OSError: pass if tmp_name is not None: try: os.unlink(tmp_name, dir_fd=directory_fd) except FileNotFoundError: pass except OSError: pass def load_config(): cfg = dict(DEFAULT_CONFIG) stored = load_json(CONFIG_FILE, {}) if isinstance(stored, dict): cfg.update(stored) for removed in ( "transcribe_whisper", "whisper_model", "max_whisper_minutes", "subs_langs", "enable_ai_summary", ): cfg.pop(removed, None) interests = cfg.get("interests", DEFAULT_CONFIG["interests"]) if not isinstance(interests, (list, tuple)): interests = DEFAULT_CONFIG["interests"] cfg["interests"] = [ keyword for keyword in _bounded_string_list(interests, 5, MAX_TAG_CHARS) if keyword.strip() ] cfg["browser"] = _clip_text(cfg.get("browser"), 32, DEFAULT_CONFIG["browser"]) cfg["feed_limit"] = _bounded_int(cfg.get("feed_limit"), 120, 1, 500) cfg["account_history_limit"] = _bounded_int( cfg.get("account_history_limit"), 200, 1, 1000 ) cfg["max_candidates"] = _bounded_int(cfg.get("max_candidates"), 60, 1, 60) cfg["metadata_workers"] = _bounded_int(cfg.get("metadata_workers"), 4, 1, 8) return cfg def save_config(cfg): ensure_dirs() atomic_write_json(CONFIG_FILE, cfg) def load_seen(): data = load_json(SEEN_FILE, {}) ids = data.get("ids", []) if isinstance(data, dict) else [] if not isinstance(ids, list): return set() return { video_id for video_id in ids[:MAX_SEEN_IDS] if isinstance(video_id, str) and VIDEO_ID_PATTERN.fullmatch(video_id) } def save_seen(seen): ensure_dirs() ids = sorted( video_id for video_id in seen if isinstance(video_id, str) and VIDEO_ID_PATTERN.fullmatch(video_id) )[-MAX_SEEN_IDS:] atomic_write_json(SEEN_FILE, {"ids": ids}) def load_dismissed(): data = load_json(DISMISSED_FILE, {}) ids = data.get("ids", []) if isinstance(data, dict) else [] if not isinstance(ids, list): ids = [] ids = [ video_id for video_id in ids[-MAX_SEEN_IDS:] if isinstance(video_id, str) and VIDEO_ID_PATTERN.fullmatch(video_id) ] return set(ids), ids def save_dismissed(dismissed_ids): """Persist dismissed video ids in dismissal order (oldest trimmed first).""" ensure_dirs() ids = [ video_id for video_id in list(dismissed_ids)[-MAX_SEEN_IDS:] if isinstance(video_id, str) and VIDEO_ID_PATTERN.fullmatch(video_id) ] atomic_write_json(DISMISSED_FILE, {"ids": ids}) RECENT_FILE = CACHE_DIR / "recent.json" RECENT_LIMIT = 5 def load_recent(): data = load_json(RECENT_FILE, {}) videos = data.get("videos", []) if isinstance(data, dict) else [] return _bounded_items(videos, RECENT_LIMIT) def remember_recent(video): """Track the last videos opened through the plugin (newest first).""" videos = [v for v in load_recent() if v.get("id") != video.get("id")] videos.insert(0, _bounded_item(video)) atomic_write_json(RECENT_FILE, {"videos": videos[:RECENT_LIMIT]}) def new_state(stage="idle", **extra): state = { "stage": stage, "detail": "", "progress": {"done": 0, "total": 0}, "interests": [], "recommendations": [], "pool": [], "recent": [], "candidates_seen": 0, "watched_count": 0, "error": None, "updated_at": time.strftime("%Y-%m-%dT%H:%M:%S"), } state.update(extra) return state def _state_payload_size(state): return len(_json_bytes(state, indent=2)) def _fit_state_size(state): """Trim optional state data until the indented persisted form fits.""" if _state_payload_size(state) <= MAX_STATE_BYTES: return state records = state["recommendations"] + state["pool"] for item in records: item["meta_description"] = item["meta_description"][:1024] item["description"] = item["description"][:320] item["tags"] = [] item["keyword_scores"] = {} if _state_payload_size(state) <= MAX_STATE_BYTES: return state # The pool is an internal ranking detail; recommendations and recent are # the only collections consumed by the UI. state["pool"] = [] if _state_payload_size(state) <= MAX_STATE_BYTES: return state recommendations = state["recommendations"] while recommendations and _state_payload_size(state) > MAX_STATE_BYTES: keep = max(0, len(recommendations) // 2) del recommendations[keep:] if _state_payload_size(state) <= MAX_STATE_BYTES: return state # This is a last-resort guard for hostile legacy state with unusually # large scalar values or record keys. state["recommendations"] = [] state["recent"] = [] return state def bound_state(state): """Sanitize state loaded from disk or received from remote metadata.""" if not isinstance(state, dict): state = {} progress = state.get("progress") if not isinstance(progress, dict): progress = {} safe = { "stage": _clip_text(state.get("stage"), 32, "idle"), "detail": _clip_text(state.get("detail"), MAX_DETAIL_CHARS), "progress": { "done": _bounded_int(progress.get("done"), 0, 0, 100000), "total": _bounded_int(progress.get("total"), 0, 0, 100000), }, "interests": [ keyword for keyword in _bounded_string_list( state.get("interests"), 5, MAX_TAG_CHARS ) if keyword.strip() ], "recommendations": _bounded_items( state.get("recommendations"), MAX_STATE_RECORDS ), "pool": _bounded_items(state.get("pool"), MAX_STATE_RECORDS), "recent": _bounded_items(state.get("recent"), RECENT_LIMIT), "candidates_seen": _bounded_int( state.get("candidates_seen"), 0, 0, 100000 ), "watched_count": _bounded_int(state.get("watched_count"), 0, 0, 1000000), "error": _clip_text(state.get("error"), MAX_ERROR_CHARS), "updated_at": _clip_text(state.get("updated_at"), 64), } return _fit_state_size(safe) def publish_state(state): """Persist and stream the current state.""" state["updated_at"] = time.strftime("%Y-%m-%dT%H:%M:%S") safe = bound_state(state) state.clear() state.update(safe) atomic_write_json(STATE_FILE, state, max_bytes=MAX_STATE_BYTES) emit_json(state, max_bytes=MAX_STATE_BYTES) # --------------------------------------------------------------------------- # watched history def video_id_from_url(url): if not isinstance(url, str) or len(url) > 4096: return None try: parsed = urlsplit(url) except ValueError: return None if parsed.scheme not in ("http", "https"): return None host = (parsed.hostname or "").lower().rstrip(".") video_id = None if host in ( "youtube.com", "www.youtube.com", "m.youtube.com", "music.youtube.com" ): if parsed.path == "/watch": values = parse_qs(parsed.query, keep_blank_values=False).get("v", []) video_id = values[0] if values else None else: parts = [part for part in parsed.path.split("/") if part] if len(parts) == 2 and parts[0] == "shorts": video_id = parts[1] elif host == "youtu.be": parts = [part for part in parsed.path.split("/") if part] if len(parts) == 1: video_id = parts[0] return video_id if isinstance(video_id, str) and VIDEO_ID_PATTERN.fullmatch(video_id) else None def collect_browser_history_ids(): ids = set() sources = [] # Firefox profiles firefox_root = HOME / ".mozilla" / "firefox" for source in _profile_db_candidates( firefox_root, "places.sqlite", MAX_HISTORY_DB_BYTES )[:MAX_BROWSER_PROFILES]: sources.append(("firefox", source)) # Chromium family browsers for name, base in ( ("chrome", HOME / ".config" / "google-chrome"), ("chromium", HOME / ".config" / "chromium"), ("brave", HOME / ".config" / "BraveSoftware" / "Brave-Browser"), ("edge", HOME / ".config" / "microsoft-edge"), ): if len(sources) >= MAX_BROWSER_PROFILES: break remaining = MAX_BROWSER_PROFILES - len(sources) for source in _profile_db_candidates( base, "History", MAX_HISTORY_DB_BYTES )[:remaining]: sources.append((name, source)) clause = " OR ".join("url LIKE ?" for _ in WATCH_PATTERNS) for name, source in sources[:MAX_BROWSER_PROFILES]: table = "moz_places" if name == "firefox" else "urls" order = "last_visit_date" if name == "firefox" else "last_visit_time" query = ( f"SELECT url FROM {table} WHERE {clause} " f"ORDER BY {order} DESC LIMIT ?" ) try: with _query_browser_snapshot(source, MAX_HISTORY_DB_BYTES) as con: params = (*WATCH_PATTERNS, MAX_HISTORY_ROWS) for (url,) in con.execute(query, params): vid = video_id_from_url(url) if vid: ids.add(vid) if len(ids) >= MAX_HISTORY_IDS: return ids except (OSError, RuntimeError, sqlite3.Error, TypeError, ValueError): continue # unreadable profile: skip silently return ids def fetch_account_history_ids(cfg, auth=None): """Best-effort: recent watch-history entries from the YouTube account. Catches videos watched on other devices; requires working cookies and history being enabled on the account. """ if auth is None: with cookie_args(cfg.get("browser")) as scoped_auth: return fetch_account_history_ids(cfg, scoped_auth) limit = _bounded_int(cfg.get("account_history_limit"), 200, 1, 1000) ids = set() try: arguments = [ "--flat-playlist", "--no-warnings", "--quiet", "--playlist-end", str(limit), "--print", "%(id)s", ":ythistory", ] proc, limited = _run_ytdlp(arguments, auth, timeout=120) if limited: return ids for line in proc.stdout.splitlines(): vid = line.strip() if VIDEO_ID_PATTERN.fullmatch(vid): ids.add(vid) if len(ids) >= limit: break except Exception: pass return ids # --------------------------------------------------------------------------- # feed fetching (yt-dlp + browser cookies) class BrowserAuth: def __init__(self, cookie_payload=None, fallback_browser=None): self.cookie_payload = cookie_payload self.fallback_browser = fallback_browser self._lock = threading.Lock() def snapshot(self): with self._lock: return self.cookie_payload def update(self, payload): with self._lock: self.cookie_payload = payload def _valid_netscape_cookie_payload(payload): if not payload or len(payload) > MAX_COOKIE_PAYLOAD_BYTES or not payload.endswith(b"\n"): return False try: lines = payload.decode("utf-8").splitlines() except UnicodeDecodeError: return False if not lines or lines[0] != "# Netscape HTTP Cookie File": return False records = 0 for line in lines[1:]: if not line or (line.startswith("#") and not line.startswith("#HttpOnly_")): continue fields = line.split("\t") if len(fields) != 7: return False host = fields[0].removeprefix("#HttpOnly_") normalized = host.lstrip(".").lower() if not any( normalized == domain or normalized.endswith("." + domain) for domain in COOKIE_DOMAINS ): return False if fields[1] not in ("TRUE", "FALSE") or fields[3] not in ("TRUE", "FALSE"): return False try: int(fields[4]) except ValueError: return False if not all(_safe_cookie_field(field) for field in (host, fields[2], fields[5], fields[6])): return False records += 1 return records > 0 def _browser_auth(browser): if not browser or browser == "none": return BrowserAuth() if browser in CHROMIUM_BROWSERS: if _AES is None: raise RuntimeError("safe Chromium cookie export requires pycryptodomex") payload = export_browser_cookies(browser) if payload is None: raise RuntimeError("could not safely export browser cookies") return BrowserAuth(cookie_payload=payload) if browser not in { "brave", "chrome", "chromium", "edge", "firefox", "opera", "vivaldi" }: raise ValueError("unsupported browser") return BrowserAuth(fallback_browser=browser) @contextmanager def cookie_args(browser): """Yield an in-memory authentication session.""" yield _browser_auth(browser) @contextmanager def _cookie_memfd(payload): if not hasattr(os, "memfd_create"): raise RuntimeError("anonymous cookie files require Linux memfd support") fd = os.memfd_create( "omarchy-youtube-cookies", getattr(os, "MFD_CLOEXEC", 0x0001) ) try: view = memoryview(payload) while view: written = os.write(fd, view) if written <= 0: raise OSError("failed to write anonymous cookie file") view = view[written:] os.lseek(fd, 0, os.SEEK_SET) yield ["--cookies", f"/proc/self/fd/{fd}"], (fd,) finally: os.close(fd) def _run_ytdlp(arguments, auth, timeout, update_auth=True): auth = auth or BrowserAuth() cookie_payload = auth.snapshot() if cookie_payload is not None: with _cookie_memfd(cookie_payload) as (auth_args, pass_fds): result = run_bounded( ["yt-dlp", *auth_args, *arguments], timeout=timeout, pass_fds=pass_fds ) process_result, limited = result if update_auth and process_result.returncode == 0 and not limited: updated = os.pread(pass_fds[0], MAX_COOKIE_PAYLOAD_BYTES + 1, 0) if _valid_netscape_cookie_payload(updated): auth.update(updated) return result auth_args = ( ["--cookies-from-browser", auth.fallback_browser] if auth.fallback_browser else [] ) return run_bounded(["yt-dlp", *auth_args, *arguments], timeout=timeout) def fetch_feed(cfg, auth=None): """Return raw entries from the logged-in subscriptions feed.""" if auth is None: with cookie_args(cfg.get("browser")) as scoped_auth: return fetch_feed(cfg, scoped_auth) feed_limit = _bounded_int(cfg.get("feed_limit"), 120, 1, 500) arguments = [ "--flat-playlist", "--no-warnings", "--quiet", "--playlist-end", str(feed_limit), "--dump-json", "https://www.youtube.com/feed/subscriptions", ] proc, limited = _run_ytdlp(arguments, auth, timeout=240) if limited: raise RuntimeError("yt-dlp feed output or timeout limit reached") entries = [] for line in proc.stdout.splitlines(): line = line.strip() if not line.startswith("{"): continue try: data = json.loads(line) except ValueError: continue if not isinstance(data, dict): continue vid = data.get("id") or "" if not isinstance(vid, str) or not VIDEO_ID_PATTERN.fullmatch(vid): continue def num(value): return _bounded_int(value, 0, 0, 10**18) entries.append( { "id": vid, "title": _clip_text( data.get("title"), MAX_TITLE_CHARS, "(untitled)" ), "channel": _clip_text( data.get("channel") or data.get("uploader"), MAX_CHANNEL_CHARS, "YouTube", ), "duration": num(data.get("duration")), "views": num(data.get("view_count")), } ) if len(entries) >= feed_limit: break if not entries and proc.returncode != 0: raise RuntimeError( "yt-dlp returned no feed entries: " + (proc.stderr.strip().splitlines() or [ "output limit reached" if limited else "unknown error" ])[-1][:200] ) return entries def fetch_recommended(cfg, limit=None, auth=None): """Return raw entries from YouTube's personalized Recommended feed. Uses the same cookie handling as the subscriptions feed. Returns [] on failure (e.g., not logged in or feed empty) so the caller can fall back to subscriptions only. """ if auth is None: with cookie_args(cfg.get("browser")) as scoped_auth: return fetch_recommended(cfg, limit, scoped_auth) feed_limit = _bounded_int(limit or cfg.get("feed_limit"), 120, 1, 500) # Both URLs are tried; yt-dlp maps them to youtube:recommended extractor. for url in ( "https://www.youtube.com/feed/recommended", ":ytrec", ): arguments = [ "--flat-playlist", "--no-warnings", "--quiet", "--playlist-end", str(feed_limit), "--dump-json", url, ] proc, limited = _run_ytdlp(arguments, auth, timeout=240) if limited: continue entries = [] for line in proc.stdout.splitlines(): line = line.strip() if not line.startswith("{"): continue try: data = json.loads(line) except ValueError: continue if not isinstance(data, dict): continue vid = data.get("id") or "" if not isinstance(vid, str) or not VIDEO_ID_PATTERN.fullmatch(vid): continue def num(value): return _bounded_int(value, 0, 0, 10**18) entries.append( { "id": vid, "title": _clip_text( data.get("title"), MAX_TITLE_CHARS, "(untitled)" ), "channel": _clip_text( data.get("channel") or data.get("uploader"), MAX_CHANNEL_CHARS, "YouTube", ), "duration": num(data.get("duration")), "views": num(data.get("view_count")), } ) if len(entries) >= feed_limit: break if entries: return entries # If first URL gave nothing but succeeded, try next URL if proc.returncode == 0: continue return [] # --------------------------------------------------------------------------- # metadata (tags/description) fetching + scoring def fetch_metadata(cfg, video_id, auth=None): """Fetch one video's metadata JSON (tags, description, …) — no download.""" if auth is None: with cookie_args(cfg.get("browser")) as scoped_auth: return fetch_metadata(cfg, video_id, scoped_auth) arguments = [ "--skip-download", "--no-warnings", "--quiet", "--dump-json", f"https://www.youtube.com/watch?v={video_id}", ] proc, limited = _run_ytdlp(arguments, auth, timeout=60, update_auth=False) if limited: return None for line in proc.stdout.splitlines(): line = line.strip() if not line.startswith("{"): continue try: data = json.loads(line) if not isinstance(data, dict): continue return { "tags": _bounded_tags(data.get("tags")), "description": _clip_text( data.get("description"), MAX_DESCRIPTION_CHARS ), "view_count": _bounded_int( data.get("view_count"), 0, 0, 10**18 ), "timestamp": _bounded_int( data.get("timestamp"), 0, 0, 10**12 ), "upload_date": _clip_text(data.get("upload_date"), 8), } except ValueError: continue return None def score_metadata(title, tags, description, keywords): """Score a video's metadata against interest keywords. Tags weigh most (explicit creator intent), then title, then description. Returns (total_score, matched_keywords, per_keyword_scores). """ total = 0 matched = [] per_keyword = {} title_lower = (title or "").lower() tags_lower = " ".join(tags or []).lower() desc_lower = (description or "").lower() for kw in keywords: kw = kw.strip() if not kw: continue pat = keyword_pattern(kw) th = len(pat.findall(title_lower)) gh = len(pat.findall(tags_lower)) dh = len(pat.findall(desc_lower)) # A single incidental description hit ("music by …" credits, sponsor # blurbs like "ai-powered dictation") is not a real match: require a # title or creator-tag hit, or several description mentions. if th + gh == 0 and dh < 2: per_keyword[kw] = 0 continue score = 3 * th + 4 * min(gh, 10) + dh per_keyword[kw] = score matched.append(kw) total += score return total, matched, per_keyword def snippet(text, keywords, max_chars=320): """First keyword-relevant part of a description, else its opening.""" if not text: return "" text = " ".join(text.split()) lower = text.lower() best = -1 for kw in keywords: kw = kw.strip() if not kw: continue m = keyword_pattern(kw).search(lower) if m and (best == -1 or m.start() < best): best = m.start() start = max(0, best - 80) if best >= 0 else 0 cut = text[start:start + max_chars] if len(text) > start + max_chars: cut = cut.rsplit(" ", 1)[0] + "…" return ("…" + cut) if start > 0 else cut # --------------------------------------------------------------------------- # scoring def keyword_pattern(kw): return re.compile(r"\b" + re.escape(kw.strip()) + r"\b", re.IGNORECASE) def format_duration(seconds): if not seconds: return "?" minutes, sec = divmod(int(seconds), 60) hours, minutes = divmod(minutes, 60) return f"{hours}:{minutes:02d}:{sec:02d}" if hours else f"{minutes}:{sec:02d}" # --------------------------------------------------------------------------- # pipeline def _base_item(entry, meta, interests): tags = _bounded_tags((meta or {}).get("tags")) description = _clip_text( (meta or {}).get("description"), MAX_DESCRIPTION_CHARS ) entry = dict(entry) entry["title"] = _clip_text(entry.get("title"), MAX_TITLE_CHARS, "(untitled)") entry["channel"] = _clip_text( entry.get("channel"), MAX_CHANNEL_CHARS, "YouTube" ) score, matched, per_keyword = score_metadata( entry["title"], tags, description, interests ) item = { **entry, "duration_formatted": format_duration(entry.get("duration")), "score": score, "matched": matched, "keyword_scores": per_keyword, "meta_description": description, "tags": tags, "thumbnail": f"https://i.ytimg.com/vi/{entry['id']}/mqdefault.jpg", "url": f"https://www.youtube.com/watch?v={entry['id']}", "description": "", } apply_trending(item, meta) return item # --------------------------------------------------------------------------- # recency + trending def _upload_timestamp(meta): """Best-effort upload unix time from yt-dlp metadata (0 when unknown).""" ts = (meta or {}).get("timestamp") if ts: try: return int(ts) except (TypeError, ValueError): pass upload_date = (meta or {}).get("upload_date") if upload_date: try: return int(time.mktime(time.strptime(upload_date, "%Y%m%d"))) except (OverflowError, ValueError): pass return 0 def apply_trending(item, meta): """Attach upload age and a trending velocity (views/hour) to an item.""" views = (meta or {}).get("view_count") or item.get("views") or 0 item["views"] = int(views or 0) ts = _upload_timestamp(meta) now = time.time() if ts > 0: age_hours = max((now - ts) / 3600.0, 0.5) item["uploaded_at"] = ts item["age_label"] = format_age(age_hours) item["views_per_hour"] = int(views / age_hours) else: item["uploaded_at"] = 0 item["age_label"] = "" item["views_per_hour"] = 0 return item def format_age(hours): if hours < 1: return f"{max(1, int(hours * 60))}m" if hours < 48: return f"{int(hours)}h" if hours < 24 * 14: return f"{int(hours / 24)}d" return f"{int(hours / (24 * 7))}w" def recency_tier(item): """0 = last 24h, 1 = last 3 days, 2 = last week, 3 = older/unknown.""" ts = item.get("uploaded_at") or 0 if not ts: return 3 age_hours = (time.time() - ts) / 3600.0 if age_hours <= 24: return 0 if age_hours <= 72: return 1 if age_hours <= 24 * 7: return 2 return 3 def pool_sort_key(item): """Latest uploads first; within the same freshness window, most trending (views/hour), then most viewed overall.""" return ( recency_tier(item), -item.get("views_per_hour", 0), -item.get("views", 0), ) def format_views_per_hour(vph): if not vph: return "" if vph >= 1000: return f"{vph / 1000:.1f}k/hr" return f"{vph}/hr" def _pick_all_for_tag(pool, tag): """All candidates that genuinely match `tag`, sorted by popularity.""" def kw_score(item): return (item.get("keyword_scores") or {}).get(tag, 0) matched = [p for p in pool if kw_score(p) > 0] matched = sorted(matched, key=pool_sort_key) out = [] for item in matched: entry = dict(item) entry["tag"] = tag # Keep snippet for list view; full description also available. entry["description"] = ( snippet(item.get("meta_description"), [tag]) or item.get("meta_description", "")[:320] or "No description available." ) out.append(entry) return out def _pick_others(pool, tags): """Videos matching no tag — sorted by popularity.""" def any_match(item): ks = item.get("keyword_scores") or {} return any(ks.get(t, 0) > 0 for t in tags) others = [p for p in pool if not any_match(p)] others = sorted(others, key=pool_sort_key) out = [] for item in others: entry = dict(item) entry["tag"] = "Others" entry["description"] = ( snippet(item.get("meta_description"), []) or item.get("meta_description", "")[:320] or "No description available." ) out.append(entry) return out def run_pipeline(limit=None, include_recommended=False): cfg = load_config() interests = [k for k in cfg.get("interests", []) if k.strip()] tags = interests # no pseudo-tag; Others handles unmatched limit = limit or cfg.get("max_candidates", 60) state = new_state(stage="history", interests=interests) publish_state(state) watched = collect_browser_history_ids() | load_seen() try: auth = _browser_auth(cfg.get("browser")) except Exception as exc: publish_state(new_state(stage="error", error=str(exc), interests=interests)) return if cfg.get("use_account_history", True): state = dict(state, detail="Fetching account watch history…") publish_state(state) watched |= fetch_account_history_ids(cfg, auth) state["watched_count"] = len(watched) state = dict(state, stage="feed", detail="Fetching subscriptions feed…") publish_state(state) try: entries = fetch_feed(cfg, auth) for e in entries: e["source"] = "subs" if include_recommended: state = dict(state, detail="Fetching YouTube Recommended…") publish_state(state) rec_entries = fetch_recommended(cfg, limit=20, auth=auth) for e in rec_entries: e["source"] = "rec" seen_ids = {e["id"] for e in entries} for e in rec_entries: if e["id"] not in seen_ids: entries.append(e) seen_ids.add(e["id"]) except Exception as exc: publish_state(new_state(stage="error", error=str(exc), interests=interests)) return if include_recommended: fresh_subs = [e for e in entries if e["id"] not in watched and e.get("source") == "subs"] fresh_recs = [e for e in entries if e["id"] not in watched and e.get("source") == "rec"] # 20 Recommended on top of subs (not interleaved within limit) candidates = fresh_subs[:limit] + fresh_recs[:20] else: fresh = [e for e in entries if e["id"] not in watched] candidates = fresh[:limit] # Score metadata (title/tags/description) concurrently; no media downloads. def _meta_job(entry): try: meta = fetch_metadata(cfg, entry["id"], auth) except Exception: meta = None return entry, meta try: pool = [None] * len(candidates) completed = 0 workers = _bounded_int(cfg.get("metadata_workers"), 4, 1, 8) with ThreadPoolExecutor(max_workers=workers) as ex: future_map = {ex.submit(_meta_job, e): i for i, e in enumerate(candidates)} for fut in as_completed(future_map): try: entry, meta = fut.result() except Exception as exc: # One fetch failed — skip it but keep pipeline alive completed += 1 state = dict( state, stage="metadata", detail=f"Skipped (fetch error: {exc})", progress={"done": completed, "total": len(candidates)}, candidates_seen=len(entries), ) publish_state(state) continue try: pool[future_map[fut]] = _base_item(entry, meta, interests) except Exception as exc: # Bad metadata should not kill the whole run pool[future_map[fut]] = { **entry, "duration_formatted": format_duration(entry.get("duration")), "score": 0, "matched": [], "keyword_scores": {}, "meta_description": _clip_text( (meta or {}).get("description"), MAX_DESCRIPTION_CHARS ), "tags": [], "thumbnail": f"https://i.ytimg.com/vi/{entry['id']}/mqdefault.jpg", "url": f"https://www.youtube.com/watch?v={entry['id']}", "description": "", "views": entry.get("views", 0), "uploaded_at": 0, "age_label": "", "views_per_hour": 0, } completed += 1 state = dict( state, stage="metadata", detail=entry["title"], progress={"done": completed, "total": len(candidates)}, candidates_seen=len(entries), ) publish_state(state) pool = sorted( [p for p in pool if p], key=pool_sort_key, ) # Keep only last 24h uploads for subs; keep all Recommended (limit 20) unfiltered now = time.time() if include_recommended: pool = [ p for p in pool if p.get("source") == "rec" or (p.get("uploaded_at") and (now - p["uploaded_at"]) < 24 * 3600) ] else: pool = [p for p in pool if p.get("uploaded_at") and (now - p["uploaded_at"]) < 24 * 3600] # Classify fast using metadata scoring; every tag gets all its matches, # Others collects everything unmatched — all sorted by popularity. recommendations = [] if tags: for tag in tags: recommendations.extend(_pick_all_for_tag(pool, tag)) recommendations.extend(_pick_others(pool, tags)) else: # No tags yet — everything goes to Others so the panel is not empty. for item in sorted(pool, key=pool_sort_key): entry = dict(item) entry["tag"] = "Others" entry["description"] = ( snippet(item.get("meta_description"), []) or item.get("meta_description", "")[:320] or "No description available." ) recommendations.append(entry) # Detail is honest when pool was empty due to 24h filter. if not pool: detail = "No uploads from subscriptions in last 24h" + (" + 0 recommended" if include_recommended else "") elif not recommendations and tags: detail = "No tag matches in last 24h — showing Others" recommendations = _pick_others(pool, tags) else: if include_recommended: rec_n = sum(1 for p in pool if p.get("source") == "rec") detail = f"{len(recommendations)} videos (subs 24h + {rec_n} recommended)" else: detail = f"{len(recommendations)} videos from last 24h" state = new_state( stage="done", detail=detail, interests=interests, recommendations=recommendations, pool=pool, recent=load_recent(), candidates_seen=len(entries), watched_count=len(watched), ) publish_state(state) except Exception as exc: import traceback err = f"{exc}\n{traceback.format_exc()[-600:]}" publish_state(new_state(stage="error", error=err[:800], interests=interests)) # --------------------------------------------------------------------------- # commands def cmd_status(_args): cfg = load_config() state = bound_state(load_json(STATE_FILE, new_state())) # Always surface the currently configured tags, even before a run. state["interests"] = cfg["interests"] state.setdefault("recent", load_recent()) state = bound_state(state) emit_json(state, max_bytes=MAX_STATE_BYTES) def cmd_watch(_args): """Stream state file changes as JSON lines (used by tooling/debugging).""" ensure_dirs() last_mtime = None emit_json(bound_state(load_json(STATE_FILE, new_state())), max_bytes=MAX_STATE_BYTES) while True: try: mtime = STATE_FILE.stat().st_mtime if mtime != last_mtime: last_mtime = mtime emit_json( bound_state(load_json(STATE_FILE, new_state())), max_bytes=MAX_STATE_BYTES, ) except OSError: pass time.sleep(1.5) def cmd_open(args): vid = args.video_id if not VIDEO_ID_PATTERN.fullmatch(vid): raise SystemExit("invalid YouTube video id") seen = load_seen() seen.add(vid) save_seen(seen) url = f"https://www.youtube.com/watch?v={vid}" state = bound_state(load_json(STATE_FILE, {})) info = {} for pool in (state.get("recommendations") or [], state.get("recent") or []): for item in pool: if isinstance(item, dict) and item.get("id") == vid: info = item break entry = { "id": vid, "title": info.get("title") or "(untitled)", "channel": info.get("channel") or "YouTube", "duration_formatted": info.get("duration_formatted") or "?", "thumbnail": f"https://i.ytimg.com/vi/{vid}/mqdefault.jpg", "url": url, "opened_at": time.strftime("%Y-%m-%dT%H:%M:%S"), } remember_recent(entry) # Keep the published state's recent list in sync for the UI. if state: state["recent"] = load_recent() atomic_write_json( STATE_FILE, bound_state(state), max_bytes=MAX_STATE_BYTES ) # Same as Super+Shift+Y: omarchy-launch-webapp does `uwsm-app -- --app=URL` try: subprocess.Popen( ["omarchy-launch-webapp", url], start_new_session=True, stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL, ) except Exception: # Fallback for headless / testing try: webbrowser.open(url) except Exception: pass emit_json({"opened": vid, "url": url}) def cmd_config(args): if args.action == "get": emit_json(load_config()) return cfg = load_config() if args.interests is not None: cfg["interests"] = [ _clip_text(k.strip(), MAX_TAG_CHARS) for k in args.interests.split(",") if k.strip() ][:5] # Sync the saved tags into the published state so the UI reflects # the change immediately (the old state file kept stale tags). state = bound_state(load_json(STATE_FILE, {})) if state: state["interests"] = cfg["interests"] atomic_write_json( STATE_FILE, bound_state(state), max_bytes=MAX_STATE_BYTES ) if args.browser: cfg["browser"] = _clip_text(args.browser, 32) if args.max_candidates: cfg["max_candidates"] = _bounded_int(args.max_candidates, 60, 1, 60) save_config(cfg) emit_json(cfg) def cmd_history(_args): ids = collect_browser_history_ids() emit_json({"watched": len(ids), "sample": sorted(ids)[:20]}) def cmd_feed(args): cfg = load_config() entries = fetch_feed(cfg) watched = collect_browser_history_ids() | load_seen() fresh = [e for e in entries if e["id"] not in watched] limit = args.limit or cfg.get("max_candidates", 15) emit_json( { "feed_size": len(entries), "unwatched": len(fresh), "candidates": fresh[:limit], } ) def main(): if len(sys.argv) >= 4 and sys.argv[1] == "_exec-child": _supervise_with_parent_death(int(sys.argv[2]), sys.argv[3:]) return parser = argparse.ArgumentParser(prog="omarchy-youtube-suggester") sub = parser.add_subparsers(dest="command", required=True) p_run = sub.add_parser("run", help="Run the full suggester pipeline") p_run.add_argument("--limit", type=int, default=None) p_run.add_argument("--with-recommended", action="store_true", dest="with_recommended", help="Also include YouTube's personalized Recommended feed (needs same cookies)") p_run.set_defaults(func=lambda a: run_pipeline(a.limit, a.with_recommended)) p_status = sub.add_parser("status", help="Print current state JSON") p_status.set_defaults(func=cmd_status) p_watch = sub.add_parser("watch", help="Stream state changes as JSON lines") p_watch.set_defaults(func=cmd_watch) p_open = sub.add_parser("open", help="Open a video and mark it watched") p_open.add_argument("video_id") p_open.set_defaults(func=cmd_open) p_cfg = sub.add_parser("config", help="Get or set configuration") p_cfg.add_argument("action", choices=["get", "set"]) p_cfg.add_argument("--interests", default=None, help="Comma-separated keywords (max 5)") p_cfg.add_argument("--browser", default=None) p_cfg.add_argument("--max-candidates", type=int, default=None, dest="max_candidates") p_cfg.set_defaults(func=cmd_config) p_hist = sub.add_parser("history", help="Rebuild watched set from browsers (debug)") p_hist.set_defaults(func=cmd_history) p_feed = sub.add_parser("feed", help="Fetch feed only (debug)") p_feed.add_argument("--limit", type=int, default=None) p_feed.set_defaults(func=cmd_feed) args = parser.parse_args() ensure_dirs() args.func(args) if __name__ == "__main__": main()