"""Scan orchestration (§4.6). Pulls libraries, items and history; normalizes; upserts inside one transaction per library; rolls episodes up to seasons; resolves keep marks; marks vanished items missing. The property that matters most here is idempotency. A scanner that double-counts sizes or duplicates history events produces a report that looks entirely plausible and is wrong, which is worse than one that crashes. """ from __future__ import annotations import json import logging import os import time from dataclasses import dataclass from . import keeps from .db import Database from .providers.base import ( Account, AuthError, Coverage, HistoryProvider, Item, Library, MediaProvider, ProviderError, WatchEvent, ) log = logging.getLogger(__name__) class ScanBusy(RuntimeError): pass @dataclass class ScanResult: scan_id: int status: str items_seen: int = 0 items_added: int = 0 items_updated: int = 0 items_missing: int = 0 events_added: int = 0 warnings: list[str] = None error: str | None = None class Ingest: def __init__(self, db: Database, cfg, media: MediaProvider, history: HistoryProvider | None): self.db = db self.cfg = cfg self.media = media self.history = history self.warnings: list[str] = [] # ── locking ────────────────────────────────────────────────────────── def _acquire_lock(self, scan_id: int) -> None: now = int(time.time()) row = self.db.one("SELECT * FROM scan_lock WHERE id = 1") if row and row["scan_id"] is not None: age = now - (row["acquired_at"] or 0) if age < self.cfg.scan_lock_timeout_s: raise ScanBusy("a scan is already running (started %ds ago)" % age) # Stale lock: the previous scan died. Mark it failed and take over. log.warning("breaking stale scan lock held by scan %s", row["scan_id"]) self.db.execute( "UPDATE scan SET status='failed', finished_at=?, " "error='abandoned - lock timed out' WHERE id=? AND status='running'", (now, row["scan_id"]), ) self.db.execute( "INSERT INTO scan_lock (id, scan_id, holder, acquired_at) VALUES (1,?,?,?) " "ON CONFLICT(id) DO UPDATE SET scan_id=excluded.scan_id, " "holder=excluded.holder, acquired_at=excluded.acquired_at", (scan_id, "%s:%s" % (os.uname().nodename, os.getpid()), now), ) def _release_lock(self) -> None: self.db.execute("UPDATE scan_lock SET scan_id=NULL, holder=NULL WHERE id=1") def _progress(self, scan_id: int, text: str) -> None: self.db.execute("UPDATE scan SET progress=? WHERE id=?", (text, scan_id)) def _warn(self, msg: str) -> None: log.warning("scan warning: %s", msg) if len(self.warnings) < 200: self.warnings.append(msg) # ── entry point ────────────────────────────────────────────────────── def run(self, mode: str = "full", trigger: str = "manual") -> ScanResult: now = int(time.time()) cur = self.db.execute( "INSERT INTO scan (mode, trigger, status, started_at, history_source) " "VALUES (?,?,'running',?,?)", (mode, trigger, now, self.history.name if self.history else None), ) scan_id = cur.lastrowid try: self._acquire_lock(scan_id) except ScanBusy: self.db.execute( "UPDATE scan SET status='failed', finished_at=?, error=? WHERE id=?", (now, "another scan is already running", scan_id), ) raise result = ScanResult(scan_id=scan_id, status="running", warnings=[]) try: self._run_inner(scan_id, mode, result) result.status = "succeeded" except AuthError as e: result.status, result.error = "failed", str(e) log.error("scan %s failed on auth: %s", scan_id, e) except Exception as e: # noqa: BLE001 result.status, result.error = "failed", str(e) log.exception("scan %s failed", scan_id) finally: result.warnings = self.warnings self.db.execute( "UPDATE scan SET status=?, finished_at=?, items_seen=?, items_added=?, " "items_updated=?, items_missing=?, events_added=?, warning_count=?, " "warnings=?, error=?, progress=NULL WHERE id=?", (result.status, int(time.time()), result.items_seen, result.items_added, result.items_updated, result.items_missing, result.events_added, len(self.warnings), json.dumps(self.warnings[:200]), result.error, scan_id), ) self._release_lock() return result def _run_inner(self, scan_id: int, mode: str, result: ScanResult) -> None: self._progress(scan_id, "connecting") info = self.media.server_info() provider_id = self._upsert_provider(info) self.db.execute("UPDATE scan SET provider_id=? WHERE id=?", (provider_id, scan_id)) self._check_history_pairing(provider_id, info) self._progress(scan_id, "reading libraries") libraries = self.media.libraries() lib_ids = self._upsert_libraries(provider_id, libraries) self._seed_keep_all_libraries() # History first: item rollups need it in place. self._progress(scan_id, "reading watch history") coverage = self._ingest_history(provider_id, mode, result) for lib in libraries: self._progress(scan_id, "scanning %s" % lib.title) self._ingest_library(provider_id, lib, lib_ids[lib.provider_key], scan_id, result) self._progress(scan_id, "rolling up") self._apply_watch_rollups(provider_id) self._correct_added_at(provider_id) self._rollup_seasons(provider_id) self._rollup_shows(provider_id) self._apply_pre_history(provider_id, coverage) if mode == "full": self._mark_missing(provider_id, scan_id, result) self._progress(scan_id, "resolving keeps") keeps.resolve_all(self.db) orphans = keeps.stamp_matches(self.db, scan_id) if orphans: self._warn("%d keep mark(s) matched nothing this scan" % orphans) self.db.execute("UPDATE provider SET last_scan_id=? WHERE id=?", (scan_id, provider_id)) # ── provider / libraries ───────────────────────────────────────────── def _upsert_provider(self, info) -> int: now = int(time.time()) self.db.execute( "INSERT INTO provider (kind, name, base_url, server_id, version, created_at) " "VALUES (?,?,?,?,?,?) ON CONFLICT(kind, base_url) DO UPDATE SET " "name=excluded.name, server_id=excluded.server_id, version=excluded.version", (info.kind, info.name, info.base_url, info.server_id, info.version, now), ) return self.db.scalar( "SELECT id FROM provider WHERE kind=? AND base_url=?", (info.kind, info.base_url), ) def _check_history_pairing(self, provider_id: int, media_info) -> None: """Refuse to join history from a different Plex server (§4.11).""" if self.history is None or self.history.name != "tautulli": return try: hinfo = self.history.server_info() except ProviderError as e: self._warn("could not read Tautulli server info: %s" % e) return if hinfo.server_id and media_info.server_id and hinfo.server_id != media_info.server_id: raise ProviderError( "Tautulli is watching a different Plex server " "(%s != %s) - refusing to join unrelated history data" % (hinfo.server_id[:8], media_info.server_id[:8]) ) def _upsert_libraries(self, provider_id: int, libraries: list[Library]) -> dict[str, int]: out = {} now = int(time.time()) for lib in libraries: self.db.execute( "INSERT INTO library (provider_id, provider_key, title, kind, locations, scanned_at) " "VALUES (?,?,?,?,?,?) ON CONFLICT(provider_id, provider_key) DO UPDATE SET " "title=excluded.title, kind=excluded.kind, locations=excluded.locations, " "scanned_at=excluded.scanned_at", (provider_id, lib.provider_key, lib.title, lib.kind, json.dumps(lib.locations), now), ) out[lib.provider_key] = self.db.scalar( "SELECT id FROM library WHERE provider_id=? AND provider_key=?", (provider_id, lib.provider_key), ) return out def _seed_keep_all_libraries(self) -> None: """Apply KEEP_ALL_LIBRARIES once, on first run only. Empty by default: nothing is ever kept unless a person says so (§6.6). Re-applying on every start would silently re-enable a rule the user turned off in the UI, so a marker setting guards it. """ if not self.cfg.keep_all_libraries: return if self.db.get_setting("keep_all_seeded"): return for title in self.cfg.keep_all_libraries: row = self.db.one("SELECT id FROM library WHERE title = ?", (title,)) if row: self.db.execute("UPDATE library SET keep_all = 1 WHERE id = ?", (row["id"],)) log.info("seeded keep_all for library %r", title) else: self._warn("KEEP_ALL_LIBRARIES names %r, which is not a library" % title) self.db.set_setting("keep_all_seeded", "1") # ── history ────────────────────────────────────────────────────────── def _disposition(self, pc: int | None) -> str: if pc is None: return "completed" # Plex fallback: only a play is recorded (§4.11) if pc >= self.cfg.completion_threshold: return "completed" if pc < self.cfg.abandon_ceiling: return "abandoned" return "partial" def _ingest_history(self, provider_id: int, mode: str, result: ScanResult) -> Coverage | None: if self.history is None: return None try: for acct in self.history.accounts(): self.db.execute( "INSERT INTO account (provider_id, account_id, name, friendly_name) " "VALUES (?,?,?,?) ON CONFLICT(provider_id, account_id) DO UPDATE SET " "name=excluded.name, friendly_name=excluded.friendly_name", (provider_id, acct.account_id, acct.name, acct.friendly_name), ) except ProviderError as e: self._warn("could not read accounts: %s" % e) since = None if mode != "full": since = self.db.scalar( "SELECT MAX(viewed_at) FROM watch_event WHERE provider_id=? AND source=?", (provider_id, self.history.name), ) window = self.cfg.session_merge_window_h * 3600 batch: list[tuple] = [] added = 0 def flush(): nonlocal added, batch if not batch: return cur = self.db.executemany( "INSERT INTO watch_event (provider_id, source, source_row_id, reference_id, " "provider_item_id, account_id, viewed_at, stopped_at, play_duration_s, " "paused_counter_s, percent_complete, watched_status, disposition, session_id, " "media_type, platform) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?) " "ON CONFLICT(provider_id, source, source_row_id) DO NOTHING", batch, ) added += cur.rowcount if cur.rowcount and cur.rowcount > 0 else 0 batch = [] for ev in self.history.watch_events(since): # Session key: same item + same user within the merge window (§4.9). bucket = ev.viewed_at // window if window > 0 else ev.viewed_at session_id = "%s:%s:%s" % (ev.provider_item_id, ev.account_id or "-", bucket) batch.append(( provider_id, ev.source, ev.source_row_id, ev.reference_id, ev.provider_item_id, ev.account_id, ev.viewed_at, ev.stopped_at, ev.play_duration_s, ev.paused_counter_s, ev.percent_complete, ev.watched_status, self._disposition(ev.percent_complete), session_id, ev.media_type, ev.platform, )) if len(batch) >= 2000: flush() flush() result.events_added = added cov = self.db.one( "SELECT MIN(viewed_at) AS lo, MAX(viewed_at) AS hi, COUNT(*) AS n " "FROM watch_event WHERE provider_id=? AND source=?", (provider_id, self.history.name), ) coverage = Coverage(cov["lo"], cov["hi"], cov["n"] or 0) self.db.execute( "INSERT INTO history_coverage (provider_id, source, earliest_event_at, " "latest_event_at, event_count, updated_at) VALUES (?,?,?,?,?,?) " "ON CONFLICT(provider_id, source) DO UPDATE SET " "earliest_event_at=excluded.earliest_event_at, " "latest_event_at=excluded.latest_event_at, " "event_count=excluded.event_count, updated_at=excluded.updated_at", (provider_id, self.history.name, coverage.earliest_event_at, coverage.latest_event_at, coverage.event_count, int(time.time())), ) return coverage # ── items ──────────────────────────────────────────────────────────── def _ingest_library(self, provider_id: int, lib: Library, library_id: int, scan_id: int, result: ScanResult) -> None: show_guids: dict[str, str] = {} if lib.kind == "show": try: show_guids = self.media.show_guids(lib) except Exception as e: # noqa: BLE001 self._warn("could not read show GUIDs for %s: %s" % (lib.title, e)) seasons: dict[str, dict] = {} shows: dict[str, dict] = {} n = 0 conn = self.db.conn conn.execute("BEGIN IMMEDIATE") try: for item in self.media.items(lib): n += 1 if item.kind == "movie": self._upsert_movie(provider_id, library_id, item, scan_id, result) else: self._collect_episode(provider_id, library_id, item, show_guids, seasons, shows, scan_id, result) # season/show container rows for key, s in seasons.items(): self._upsert_season(provider_id, library_id, s, scan_id, result) for key, s in shows.items(): self._upsert_show(provider_id, library_id, s, scan_id, result) conn.execute("COMMIT") except Exception: conn.execute("ROLLBACK") raise result.items_seen += n def _upsert_movie(self, provider_id, library_id, item: Item, scan_id, result) -> None: existing = self.db.one( "SELECT id FROM media_item WHERE provider_id=? AND provider_item_id=?", (provider_id, item.provider_item_id), ) primary = item.parts[0].file_path if item.parts else None now = int(time.time()) vals = ( provider_id, library_id, item.provider_item_id, "movie", item.guid, None, item.title, item.sort_title, item.year, None, None, item.added_at, item.added_at, item.updated_at, 0, item.size_bytes, item.duration_ms, len(item.parts), primary, item.resolution, item.video_codec, item.view_count, item.last_viewed_at, "present", scan_id, scan_id, now, ) if existing: self.db.execute( "UPDATE media_item SET library_id=?, guid=?, title=?, sort_title=?, year=?, " "provider_added_at=?, updated_at=?, size_bytes=?, duration_ms=?, part_count=?, " "primary_path=?, resolution=?, video_codec=?, provider_view_count=?, " "provider_last_viewed_at=?, status='present', last_seen_scan_id=? WHERE id=?", (library_id, item.guid, item.title, item.sort_title, item.year, item.added_at, item.updated_at, item.size_bytes, item.duration_ms, len(item.parts), primary, item.resolution, item.video_codec, item.view_count, item.last_viewed_at, scan_id, existing["id"]), ) item_id = existing["id"] result.items_updated += 1 else: cur = self.db.execute( "INSERT INTO media_item (provider_id, library_id, provider_item_id, kind, " "guid, show_guid, title, sort_title, year, parent_id, season_number, " "added_at, provider_added_at, updated_at, episode_count, size_bytes, " "duration_ms, part_count, primary_path, resolution, video_codec, " "provider_view_count, provider_last_viewed_at, status, " "first_seen_scan_id, last_seen_scan_id, first_seen_at) " "VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)", vals, ) item_id = cur.lastrowid result.items_added += 1 # Parts are replaced wholesale — cheap, and the only way to stay correct # when a version is removed or a split file is re-encoded into one. self.db.execute("DELETE FROM media_part WHERE media_item_id=?", (item_id,)) self._insert_parts(item, media_item_id=item_id) def _insert_parts(self, item: Item, *, media_item_id=None, episode_id=None) -> None: if not item.parts: return self.db.executemany( "INSERT INTO media_part (media_item_id, episode_id, provider_part_id, file_path, " "size_bytes, container, resolution, video_codec, audio_codec, bitrate) " "VALUES (?,?,?,?,?,?,?,?,?,?)", [(media_item_id, episode_id, p.provider_part_id, p.file_path, p.size_bytes, p.container, p.resolution, p.video_codec, p.audio_codec, p.bitrate) for p in item.parts], ) def _collect_episode(self, provider_id, library_id, item: Item, show_guids, seasons, shows, scan_id, result) -> None: if not item.season_id: self._warn("episode %s has no season; skipped" % item.provider_item_id) return show_guid = show_guids.get(item.show_id or "") or None s = seasons.setdefault(item.season_id, { "provider_item_id": item.season_id, "show_id": item.show_id, "show_guid": show_guid, "show_title": item.show_title, "season_number": item.season_number, "episodes": [], }) s["episodes"].append(item) sh = shows.setdefault(item.show_id or "?", { "provider_item_id": item.show_id, "guid": show_guid, "title": item.show_title or "(unknown show)", "seasons": set(), }) sh["seasons"].add(item.season_id) def _upsert_season(self, provider_id, library_id, s: dict, scan_id, result) -> None: eps: list[Item] = s["episodes"] title = "Season %s" % (s["season_number"] if s["season_number"] is not None else "?") existing = self.db.one( "SELECT id FROM media_item WHERE provider_id=? AND provider_item_id=?", (provider_id, s["provider_item_id"]), ) if existing: season_id = existing["id"] self.db.execute( "UPDATE media_item SET library_id=?, show_guid=?, title=?, season_number=?, " "status='present', last_seen_scan_id=? WHERE id=?", (library_id, s["show_guid"], title, s["season_number"], scan_id, season_id), ) result.items_updated += 1 else: cur = self.db.execute( "INSERT INTO media_item (provider_id, library_id, provider_item_id, kind, " "show_guid, title, season_number, status, first_seen_scan_id, last_seen_scan_id) " "VALUES (?,?,?,'season',?,?,?,'present',?,?)", (provider_id, library_id, s["provider_item_id"], s["show_guid"], title, s["season_number"], scan_id, scan_id), ) season_id = cur.lastrowid result.items_added += 1 for ep in eps: row = self.db.one("SELECT id FROM episode WHERE provider_item_id=?", (ep.provider_item_id,)) primary_count = len(ep.parts) if row: ep_id = row["id"] self.db.execute( "UPDATE episode SET season_item_id=?, episode_number=?, title=?, " "provider_added_at=?, provider_last_viewed_at=?, duration_ms=?, " "size_bytes=?, part_count=?, " "status='present', last_seen_scan_id=? WHERE id=?", (season_id, ep.episode_number, ep.title, ep.added_at, ep.last_viewed_at, ep.duration_ms, ep.size_bytes, primary_count, scan_id, ep_id), ) else: cur = self.db.execute( "INSERT INTO episode (season_item_id, provider_item_id, episode_number, " "title, added_at, provider_added_at, provider_last_viewed_at, " "duration_ms, size_bytes, part_count, " "status, last_seen_scan_id, first_seen_at) " "VALUES (?,?,?,?,?,?,?,?,?,?, 'present', ?, ?)", (season_id, ep.provider_item_id, ep.episode_number, ep.title, ep.added_at, ep.added_at, ep.last_viewed_at, ep.duration_ms, ep.size_bytes, primary_count, scan_id, int(time.time())), ) ep_id = cur.lastrowid self.db.execute("DELETE FROM media_part WHERE episode_id=?", (ep_id,)) self._insert_parts(ep, episode_id=ep_id) def _upsert_show(self, provider_id, library_id, sh: dict, scan_id, result) -> None: if not sh["provider_item_id"]: return existing = self.db.one( "SELECT id FROM media_item WHERE provider_id=? AND provider_item_id=?", (provider_id, sh["provider_item_id"]), ) if existing: show_id = existing["id"] self.db.execute( "UPDATE media_item SET library_id=?, guid=?, title=?, status='present', " "last_seen_scan_id=? WHERE id=?", (library_id, sh["guid"], sh["title"], scan_id, show_id), ) else: cur = self.db.execute( "INSERT INTO media_item (provider_id, library_id, provider_item_id, kind, " "guid, title, status, first_seen_scan_id, last_seen_scan_id) " "VALUES (?,?,?,'show',?,?,'present',?,?)", (provider_id, library_id, sh["provider_item_id"], sh["guid"], sh["title"], scan_id, scan_id), ) show_id = cur.lastrowid # link seasons to their show self.db.execute( "UPDATE media_item SET parent_id=? WHERE provider_id=? AND kind='season' " "AND provider_item_id IN (%s)" % ",".join("?" * len(sh["seasons"])), tuple([show_id, provider_id] + list(sh["seasons"])), ) # ── rollups ────────────────────────────────────────────────────────── def _apply_watch_rollups(self, provider_id: int) -> None: """Aggregate watch_event onto movies and episodes. Distinct sessions, not raw events: a paused-and-resumed play is one viewing (§4.9). Everything is recomputed from scratch each scan, which is what makes re-running a scan idempotent. """ c = self.db.conn c.execute(""" UPDATE media_item SET watch_count=0, partial_count=0, abandoned_count=0, last_watched_at=NULL, last_touched_at=NULL, first_watched_at=NULL, distinct_watcher_count=0, avg_percent_complete=NULL WHERE kind='movie' """) c.execute(""" UPDATE episode SET watch_count=0, partial_count=0, abandoned_count=0, last_watched_at=NULL, last_touched_at=NULL, first_watched_at=NULL """) agg = """ SELECT provider_item_id AS pid, COUNT(DISTINCT CASE WHEN disposition='completed' THEN session_id END) AS completed, COUNT(DISTINCT CASE WHEN disposition='partial' THEN session_id END) AS partial, COUNT(DISTINCT CASE WHEN disposition='abandoned' THEN session_id END) AS abandoned, MAX(CASE WHEN disposition='completed' THEN viewed_at END) AS last_watched, MIN(CASE WHEN disposition='completed' THEN viewed_at END) AS first_watched, MAX(viewed_at) AS last_touched, COUNT(DISTINCT account_id) AS watchers, AVG(percent_complete) AS avg_pc FROM watch_event WHERE provider_id = ? GROUP BY provider_item_id """ c.execute("DROP TABLE IF EXISTS _wagg") c.execute("CREATE TEMP TABLE _wagg AS " + agg, (provider_id,)) c.execute("CREATE INDEX _wagg_pid ON _wagg(pid)") c.execute(""" UPDATE media_item SET watch_count = COALESCE((SELECT completed FROM _wagg WHERE pid = media_item.provider_item_id), 0), partial_count = COALESCE((SELECT partial FROM _wagg WHERE pid = media_item.provider_item_id), 0), abandoned_count = COALESCE((SELECT abandoned FROM _wagg WHERE pid = media_item.provider_item_id), 0), last_watched_at = (SELECT last_watched FROM _wagg WHERE pid = media_item.provider_item_id), first_watched_at = (SELECT first_watched FROM _wagg WHERE pid = media_item.provider_item_id), last_touched_at = (SELECT last_touched FROM _wagg WHERE pid = media_item.provider_item_id), distinct_watcher_count = COALESCE((SELECT watchers FROM _wagg WHERE pid = media_item.provider_item_id), 0), avg_percent_complete = (SELECT avg_pc FROM _wagg WHERE pid = media_item.provider_item_id) WHERE kind = 'movie' """) c.execute(""" UPDATE episode SET watch_count = COALESCE((SELECT completed FROM _wagg WHERE pid = episode.provider_item_id), 0), partial_count = COALESCE((SELECT partial FROM _wagg WHERE pid = episode.provider_item_id), 0), abandoned_count = COALESCE((SELECT abandoned FROM _wagg WHERE pid = episode.provider_item_id), 0), last_watched_at = (SELECT last_watched FROM _wagg WHERE pid = episode.provider_item_id), first_watched_at = (SELECT first_watched FROM _wagg WHERE pid = episode.provider_item_id), last_touched_at = (SELECT last_touched FROM _wagg WHERE pid = episode.provider_item_id) """) def _correct_added_at(self, provider_id: int) -> None: """Repair added_at where Plex's value is provably wrong. Plex's addedAt tracks the FILE, not the library entry: replace or re-encode a file and Date Added resets while the item, its ratingKey and its watch history all survive. Measured on the live server, 19% of items report a view EARLIER than their added date. A view is proof the item already existed, so the earliest view is a lower bound on the true add date. Two independent witnesses: * Plex's own lastViewedAt, which reaches back as far as the server does * MediaShelf's first completed play from Tautulli, which is more precise but only covers the history window (here, 2025-03 onward) The first version of this used only Tautulli and fired zero times, because the wrong dates are mostly older than the history window. Plex's lastViewedAt is the field that actually carries the evidence. This is a lower bound, not the date Plex never kept. But it beats a value we can prove impossible, and it matters: pre_history is derived from added_at, so a wrongly-recent date promotes an item into the CONFIDENT reclaim pool when it belongs in the uncertain one. Items nobody has ever watched keep Plex's value; nothing contradicts it. """ c = self.db.conn # episodes first, so season rollups inherit corrected dates c.execute("UPDATE episode SET added_at = provider_added_at " "WHERE provider_added_at IS NOT NULL") c.execute(""" UPDATE episode SET added_at = MIN( COALESCE(added_at, 253402300799), COALESCE(NULLIF(first_watched_at, 0), 253402300799), COALESCE(NULLIF(provider_last_viewed_at, 0), 253402300799)) WHERE (first_watched_at > 0 OR provider_last_viewed_at > 0) """) c.execute("UPDATE media_item SET added_at = provider_added_at, " "added_at_source = 'provider' " "WHERE kind = 'movie' AND provider_added_at IS NOT NULL " "AND provider_id = ?", (provider_id,)) cur = c.execute(""" UPDATE media_item SET added_at = MIN( COALESCE(added_at, 253402300799), COALESCE(NULLIF(first_watched_at, 0), 253402300799), COALESCE(NULLIF(provider_last_viewed_at, 0), 253402300799)), added_at_source = 'first_watch' WHERE kind = 'movie' AND provider_id = ? AND (first_watched_at > 0 OR provider_last_viewed_at > 0) AND MIN(COALESCE(NULLIF(first_watched_at, 0), 253402300799), COALESCE(NULLIF(provider_last_viewed_at, 0), 253402300799)) < COALESCE(added_at, 253402300799) """, (provider_id,)) corrected = cur.rowcount if cur.rowcount and cur.rowcount > 0 else 0 if corrected: self._warn( "%d movie(s) had a Plex addedAt later than a recorded view; " "corrected to the earliest view (Plex's Date Added follows the " "file, so replacing one resets it)" % corrected) def _rollup_seasons(self, provider_id: int) -> None: c = self.db.conn c.execute(""" UPDATE media_item SET episode_count = COALESCE((SELECT COUNT(*) FROM episode e WHERE e.season_item_id = media_item.id AND e.status='present'), 0), size_bytes = COALESCE((SELECT SUM(e.size_bytes) FROM episode e WHERE e.season_item_id = media_item.id AND e.status='present'), 0), duration_ms = COALESCE((SELECT SUM(e.duration_ms) FROM episode e WHERE e.season_item_id = media_item.id AND e.status='present'), 0), part_count = COALESCE((SELECT SUM(e.part_count) FROM episode e WHERE e.season_item_id = media_item.id AND e.status='present'), 0), added_at = (SELECT MIN(e.added_at) FROM episode e WHERE e.season_item_id = media_item.id AND e.status='present'), watch_count = COALESCE((SELECT SUM(e.watch_count) FROM episode e WHERE e.season_item_id = media_item.id), 0), partial_count = COALESCE((SELECT SUM(e.partial_count) FROM episode e WHERE e.season_item_id = media_item.id), 0), abandoned_count = COALESCE((SELECT SUM(e.abandoned_count) FROM episode e WHERE e.season_item_id = media_item.id), 0), last_watched_at = (SELECT MAX(e.last_watched_at) FROM episode e WHERE e.season_item_id = media_item.id), last_touched_at = (SELECT MAX(e.last_touched_at) FROM episode e WHERE e.season_item_id = media_item.id) WHERE kind = 'season' AND provider_id = ? """, (provider_id,)) c.execute(""" UPDATE media_item SET provider_added_at = ( SELECT MIN(e.provider_added_at) FROM episode e WHERE e.season_item_id = media_item.id AND e.status='present') WHERE kind = 'season' AND provider_id = ? """, (provider_id,)) c.execute(""" UPDATE media_item SET added_at_source = CASE WHEN provider_added_at IS NOT NULL AND added_at < provider_added_at THEN 'first_watch' ELSE 'provider' END WHERE kind = 'season' AND provider_id = ? """, (provider_id,)) # distinct watchers across the season's episodes c.execute(""" UPDATE media_item SET distinct_watcher_count = COALESCE(( SELECT COUNT(DISTINCT w.account_id) FROM watch_event w JOIN episode e ON e.provider_item_id = w.provider_item_id WHERE e.season_item_id = media_item.id), 0) WHERE kind = 'season' AND provider_id = ? """, (provider_id,)) # representative path: the common directory of its episodes c.execute(""" UPDATE media_item SET primary_path = ( SELECT p.file_path FROM media_part p JOIN episode e ON e.id = p.episode_id WHERE e.season_item_id = media_item.id ORDER BY e.episode_number LIMIT 1) WHERE kind = 'season' AND provider_id = ? """, (provider_id,)) c.execute(""" UPDATE media_item SET resolution = ( SELECT p.resolution FROM media_part p JOIN episode e ON e.id = p.episode_id WHERE e.season_item_id = media_item.id AND p.resolution IS NOT NULL LIMIT 1) WHERE kind = 'season' AND provider_id = ? """, (provider_id,)) def _rollup_shows(self, provider_id: int) -> None: self.db.execute(""" UPDATE media_item SET episode_count = COALESCE((SELECT SUM(s.episode_count) FROM media_item s WHERE s.parent_id = media_item.id), 0), size_bytes = COALESCE((SELECT SUM(s.size_bytes) FROM media_item s WHERE s.parent_id = media_item.id), 0), part_count = COALESCE((SELECT SUM(s.part_count) FROM media_item s WHERE s.parent_id = media_item.id), 0), duration_ms = COALESCE((SELECT SUM(s.duration_ms) FROM media_item s WHERE s.parent_id = media_item.id), 0), added_at = (SELECT MIN(s.added_at) FROM media_item s WHERE s.parent_id = media_item.id), watch_count = COALESCE((SELECT SUM(s.watch_count) FROM media_item s WHERE s.parent_id = media_item.id), 0), abandoned_count = COALESCE((SELECT SUM(s.abandoned_count) FROM media_item s WHERE s.parent_id = media_item.id), 0), last_watched_at = (SELECT MAX(s.last_watched_at) FROM media_item s WHERE s.parent_id = media_item.id), last_touched_at = (SELECT MAX(s.last_touched_at) FROM media_item s WHERE s.parent_id = media_item.id) WHERE kind = 'show' AND provider_id = ? """, (provider_id,)) def _apply_pre_history(self, provider_id: int, coverage: Coverage | None) -> None: """Flag items added before watch history began (§4.11). On the measured library this is the majority state, not an edge case. """ self.db.execute("UPDATE media_item SET pre_history = 0 WHERE provider_id = ?", (provider_id,)) if not coverage or not coverage.earliest_event_at: return self.db.execute( "UPDATE media_item SET pre_history = 1 " "WHERE provider_id = ? AND added_at IS NOT NULL AND added_at > 0 AND added_at < ?", (provider_id, coverage.earliest_event_at), ) def _mark_missing(self, provider_id: int, scan_id: int, result: ScanResult) -> None: """Items not seen this full sweep become 'missing', never deleted (§5.4).""" cur = self.db.execute( "UPDATE media_item SET status='missing' " "WHERE provider_id=? AND status='present' " "AND (last_seen_scan_id IS NULL OR last_seen_scan_id != ?)", (provider_id, scan_id), ) result.items_missing = cur.rowcount if cur.rowcount and cur.rowcount > 0 else 0 self.db.execute( "UPDATE episode SET status='missing' " "WHERE status='present' AND (last_seen_scan_id IS NULL OR last_seen_scan_id != ?)", (scan_id,), )