The application the design describes: Flask + SQLite, Plex for library data, Tautulli for watch history, report-only. Structure follows the design's seams. providers/ splits MediaProvider from HistoryProvider, because on this network library data and watch data live on different machines and Jellyfin later will have no Tautulli equivalent. scoring.py implements the reclaim score twice - as a SQL expression for the live grid (weights change on every slider drag, so storing it would mean rewriting thousands of rows per drag) and in Python for CSV export and tests, with a property test over 500 generated rows asserting the two agree. rules.py compiles saved views to parameterized SQL through a field/operator whitelist; nothing user-supplied is ever interpolated. Three properties are enforced by test rather than asserted in prose: - Ingest is idempotent. Three consecutive full scans leave every count and every byte total unchanged. A scanner that double-counts produces a report that looks plausible and is wrong. - Keep marks survive Plex reassigning every rating key in the library. They are keyed on content GUID, scoped per library so the Movies and 4K Movies copies of the same film mark independently. - Every config variable the app reads is declared in docker-compose.yml, so a variable set in Portainer can never silently do nothing. Also found and fixed while verifying against a fake Plex+Tautulli pair: executescript() commits the pending transaction, so migrations needed their BEGIN/COMMIT inside the script; replaceChildren() renders null as the literal text "null"; a hash-only URL change does not reload the document, so deep links needed a hashchange listener; and SQLite ROUND rounds half away from zero where Python rounds half to even. 73 tests, no live server required. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01GVbG48GAXfCZatcmX123Ra
696 lines
33 KiB
Python
696 lines
33 KiB
Python
"""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._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
|
|
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.updated_at, 0, item.size_bytes, item.duration_ms,
|
|
len(item.parts), primary, item.resolution, item.video_codec,
|
|
item.view_count, "present", scan_id, scan_id,
|
|
)
|
|
if existing:
|
|
self.db.execute(
|
|
"UPDATE media_item SET library_id=?, guid=?, title=?, sort_title=?, year=?, "
|
|
"added_at=?, updated_at=?, size_bytes=?, duration_ms=?, part_count=?, "
|
|
"primary_path=?, resolution=?, video_codec=?, provider_view_count=?, "
|
|
"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, 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, updated_at, episode_count, size_bytes, duration_ms, part_count, "
|
|
"primary_path, resolution, video_codec, provider_view_count, status, "
|
|
"first_seen_scan_id, last_seen_scan_id) "
|
|
"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=?, "
|
|
"added_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.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, duration_ms, size_bytes, part_count, status, last_seen_scan_id) "
|
|
"VALUES (?,?,?,?,?,?,?,?, 'present', ?)",
|
|
(season_id, ep.provider_item_id, ep.episode_number, ep.title,
|
|
ep.added_at, ep.duration_ms, ep.size_bytes, primary_count, scan_id),
|
|
)
|
|
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
|
|
""")
|
|
|
|
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),
|
|
last_touched_at = (SELECT last_touched FROM _wagg WHERE pid = episode.provider_item_id)
|
|
""")
|
|
|
|
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,))
|
|
|
|
# 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,),
|
|
)
|