MediaShelf/mediashelf/ingest.py
Jess Hallsworth e819548ff2
Stop trusting Plex's Date Added on its own
Jess spotted that dates looked like file dates rather than library-add dates.
He is right, and MediaShelf was not the culprit: it reproduces Plex's addedAt
exactly (verified 500/500 identical to the second). Plex's own field is what
follows the file — replace or re-encode one and Date Added resets while the
item, its ratingKey and its watch history all survive.

Measured on the live library, comparing addedAt against lastViewedAt where both
exist: 55 of 509 movies (10.8%) and 306 of 1,393 TV Show Archive items (22.0%)
were watched BEFORE they were "added" — 19% overall. 2001: A Space Odyssey
reports added 2026-07-31, last watched 2017-08-26.

That is not cosmetic. pre_history is derived from added_at, so an old item whose
file was replaced looks post-coverage and gets promoted into the CONFIDENT
reclaim pool, which is the one pool meant to be trustworthy.

A completed play proves the item already existed, so added_at is now
MIN(provider_added_at, first_watched_at). Plex's raw value is kept in
provider_added_at, added_at_source records which applied, and the item drawer
explains the substitution instead of quietly disagreeing with Plex. Unwatched
items keep Plex's value since nothing contradicts it. first_seen_at is also
recorded now and is authoritative for anything added from here on.

Plex's API has no better field; the true insert time is only in Plex's own
metadata_items.created_at on Loki.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01GVbG48GAXfCZatcmX123Ra
2026-09-10 17:35:48 +00:00

753 lines
36 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._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, "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=?, "
"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, provider_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, "
"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=?, 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, provider_added_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.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 and its watch
history survive. Measured on the live server, 19% of items report a
last-watch EARLIER than their added date.
A play is proof the item already existed, so the first completed view is
a lower bound on the true add date. That is not the exact date Plex
never kept, but it is strictly better than a value we can show is
impossible — and it matters, because pre_history is derived from
added_at and a wrongly-recent date promotes an item into the CONFIDENT
reclaim pool when it belongs in the uncertain one.
Items nobody ever watched keep Plex's value; nothing contradicts it.
first_seen_at is authoritative for anything MediaShelf sees appear.
"""
c = self.db.conn
c.execute("UPDATE episode SET added_at = provider_added_at "
"WHERE provider_added_at IS NOT NULL")
c.execute("""
UPDATE episode SET added_at = first_watched_at
WHERE first_watched_at IS NOT NULL AND first_watched_at > 0
AND (added_at IS NULL OR added_at > first_watched_at)
""")
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 = first_watched_at,
added_at_source = 'first_watch'
WHERE kind = 'movie' AND provider_id = ?
AND first_watched_at IS NOT NULL AND first_watched_at > 0
AND (added_at IS NULL OR added_at > first_watched_at)
""", (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 their first play; "
"corrected to the first play date" % 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 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,),
)