Files
Magent/backend/app/services/insights.py
T
Assclaw 333a799e21
Magent CI/CD / verify (push) Successful in 10m58s
Magent CI/CD / deploy-prod (push) Skipped
Magent CI/CD / deploy-beta (push) Successful in 1m42s
Add personal monthly viewing reports and CSV exports
2026-09-09 16:25:35 +12:00

233 lines
12 KiB
Python

import asyncio
import hashlib
import json
import math
import sqlite3
import time
from collections import defaultdict
from contextlib import closing
from datetime import datetime, timedelta, timezone
from .. import db
from ..clients.jellyfin import JellyfinClient
from ..clients.jellystat import JellystatClient, JellystatError
from ..runtime import get_runtime_settings
from .jellyfin_identity import link_user, linked_user_id
from .insights_artwork import item_id as artwork_item_id, with_artwork
_cache: dict[tuple, tuple[float, dict]] = {}
CACHE_SECONDS = 60
HARDWARE = {"amf": "AMD AMF", "qsv": "Intel Quick Sync", "nvenc": "NVIDIA NVENC",
"v4l2m2m": "V4L2", "vaapi": "VAAPI", "videotoolbox": "Apple VideoToolbox", "rkmpp": "Rockchip MPP"}
HARDWARE_ENUM = {0: "none", 1: "amf", 2: "qsv", 3: "nvenc", 4: "v4l2m2m", 5: "vaapi", 6: "videotoolbox", 7: "rkmpp"}
def add_transcoding(row, duration, media_type, totals, hardware, audio_codecs):
# Jellystat can retain stale transcoding metadata after a switch to DirectPlay.
method = row.get("PlayMethod")
if method not in {"Transcode", "DirectStream"}:
return
info = row.get("TranscodingInfo")
if isinstance(info, str):
try:
info = json.loads(info)
except ValueError:
info = None
info = info if isinstance(info, dict) else {}
video_present = media_type in {"movie", "episode"} or bool(info.get("VideoCodec"))
if method == "Transcode" and video_present:
if info.get("IsVideoDirect") is False:
totals["video_minutes"] += duration
value = info.get("HardwareAccelerationType")
value = HARDWARE_ENUM.get(value) if type(value) is int else str(value or "").strip().lower()
if value in HARDWARE:
totals["hardware_video_minutes"] += duration
hardware[HARDWARE[value]] += duration
elif value == "none":
totals["software_video_minutes"] += duration
else:
totals["unknown_hardware_minutes"] += duration
elif info.get("IsVideoDirect") is not True:
totals["unknown_video_minutes"] += duration
if info.get("IsAudioDirect") is False:
totals["audio_minutes"] += duration
codec = str(info.get("AudioCodec") or "Unknown").upper()[:30]
audio_codecs[codec] += duration
elif info.get("IsAudioDirect") is not True:
totals["unknown_audio_minutes"] += duration
def _date(value) -> datetime:
try:
result = datetime.fromisoformat(str(value).replace("Z", "+00:00"))
return result.replace(tzinfo=timezone.utc) if result.tzinfo is None else result.astimezone(timezone.utc)
except (ValueError, TypeError) as exc:
raise JellystatError("Jellystat returned an invalid history date") from exc
def _duration(value) -> float:
try:
result = float(value or 0)
if not math.isfinite(result) or result < 0:
raise ValueError()
return result
except (ValueError, TypeError, OverflowError) as exc:
raise JellystatError("Jellystat returned an invalid playback duration") from exc
async def resolve_identity(user: dict, runtime) -> str | None:
identity = await asyncio.to_thread(linked_user_id, user["username"], runtime.jellyfin_base_url)
if identity:
return identity
if user.get("auth_provider") != "jellyfin":
return None
# Bootstrap existing Jellyfin accounts from the canonical server, using exact names.
# Local accounts and email-prefix matches cannot claim a Jellyfin identity.
client = JellyfinClient(runtime.jellyfin_base_url, runtime.jellyfin_api_key)
if not client.configured():
return None
try:
users = await client.get_users()
except Exception as exc:
raise JellystatError("Could not resolve the linked Jellyfin account") from exc
matches = [entry for entry in users if isinstance(entry, dict)
and str(entry.get("Name") or "").strip().casefold() == user["username"].strip().casefold()] if isinstance(users, list) else []
if len(matches) != 1 or not matches[0].get("Id"):
return None
await asyncio.to_thread(link_user, user["username"], str(matches[0]["Id"]), runtime.jellyfin_base_url)
return await asyncio.to_thread(linked_user_id, user["username"], runtime.jellyfin_base_url)
def request_summary(user: dict, start: datetime, end: datetime, *, end_exclusive: bool = False) -> dict:
operator = "<" if end_exclusive else "<="
clause = f"julianday(created_at) >= julianday(?) AND julianday(created_at) {operator} julianday(?)"
params = [start.isoformat(), end.isoformat()]
if user.get("jellyseerr_user_id") is not None:
clause += " AND requested_by_id = ?"
params.append(user["jellyseerr_user_id"])
else:
clause += " AND requested_by_id IS NULL AND lower(trim(requested_by)) = ?"
params.append(user["username"].strip().lower())
with closing(db._connect()) as conn, conn:
conn.row_factory = sqlite3.Row
counts = conn.execute(f"""SELECT COUNT(*) AS total,
COALESCE(SUM(media_type = 'movie'), 0) AS movies,
COALESCE(SUM(media_type = 'tv'), 0) AS tv,
COALESCE(SUM(status = 1), 0) AS pending,
COALESCE(SUM(status = 2), 0) AS approved,
COALESCE(SUM(status = 3), 0) AS declined FROM requests_cache WHERE {clause}""", params).fetchone()
recent = conn.execute(f"""SELECT request_id, title, media_type, status FROM requests_cache
WHERE {clause} ORDER BY created_at DESC LIMIT 5""", params).fetchall()
return {**dict(counts), "recent": [dict(row) for row in recent]}
def summarize(history: list, libraries: list, start: datetime, end: datetime, *, end_exclusive: bool = False) -> dict:
library_types = {str(row.get("Id")): str(row.get("CollectionType") or "").lower() for row in libraries}
daily_seconds = defaultdict(float)
clients = defaultdict(float)
methods = defaultdict(float)
transcoding = dict.fromkeys(("video_minutes", "audio_minutes", "hardware_video_minutes", "software_video_minutes",
"unknown_hardware_minutes", "unknown_video_minutes", "unknown_audio_minutes"), 0.0)
hardware, audio_codecs = defaultdict(float), defaultdict(float)
titles = {}
movie_ids, episode_ids, seen = set(), set(), set()
recent = []
seconds = 0.0
for row in history:
row_id = str(row.get("Id") or "")
if not row_id:
raise JellystatError("Jellystat returned history without an activity ID")
if row_id in seen:
continue
seen.add(row_id)
date = _date(row.get("ActivityDateInserted"))
# Defend against older upstream versions ignoring the range filter.
if date < start or (date >= end if end_exclusive else date > end):
continue
duration = _duration(row.get("PlaybackDuration"))
if duration <= 0:
continue
item_id = str(row.get("NowPlayingItemId") or row_id)
episode_id = row.get("EpisodeId")
library_type = library_types.get(str(row.get("ParentId")), "")
media_type = "episode" if episode_id else "movie" if library_type == "movies" else "other"
add_transcoding(row, duration / 60, media_type, transcoding, hardware, audio_codecs)
if media_type == "episode":
episode_ids.add(str(episode_id))
elif media_type == "movie":
movie_ids.add(item_id)
seconds += duration
daily_seconds[date.date().isoformat()] += duration
client = str(row.get("Client") or "Unknown player")[:200]
clients[client] += duration
method = str(row.get("PlayMethod") or "Unknown")
method = {"DirectPlay": "Direct play", "DirectStream": "Direct stream", "Transcode": "Transcode"}.get(method, "Other")
methods[method] += duration
name = str(row.get("NowPlayingItemName") or "Untitled")[:500]
series = str(row.get("SeriesName") or "")[:500]
title = titles.setdefault(item_id, {"title": series or name, "type": "series" if episode_id else media_type, "minutes": 0, "plays": 0})
title["minutes"] += duration / 60
title["plays"] += 1
recent.append({"id": row_id, "title": name, "series": series, "type": media_type,
"episode": f"S{row.get('SeasonNumber', '?')} · E{row.get('EpisodeNumber', '?')}" if episode_id else None,
"minutes": round(duration / 60, 1), "played_at": date.isoformat(), "client": client,
"method": method, "artwork_item_id": artwork_item_id(row.get("NowPlayingItemId"))})
last_date = (end - timedelta(microseconds=1)).date() if end_exclusive and end > start else end.date()
count = (last_date - start.date()).days + 1
daily = [{"date": (start.date() + timedelta(days=i)).isoformat(),
"minutes": round(daily_seconds.get((start.date() + timedelta(days=i)).isoformat(), 0) / 60, 2)} for i in range(count)]
active_days = {day for day, duration in daily_seconds.items() if duration >= 60}
longest = run = 0
for day in daily:
run = run + 1 if day["date"] in active_days else 0
longest = max(longest, run)
current = 0
cursor = last_date if last_date.isoformat() in active_days else last_date - timedelta(days=1)
while cursor.isoformat() in active_days:
current += 1
cursor -= timedelta(days=1)
top = sorted(titles.values(), key=lambda row: (-row["minutes"], row["title"]))[:6]
for row in top:
row["minutes"] = round(row["minutes"], 1)
return {"summary": {"minutes": round(seconds / 60, 1), "plays": len(recent), "movies": len(movie_ids),
"episodes": len(episode_ids), "active_days": len(active_days),
"current_streak": current, "longest_streak": longest},
"daily": daily, "top_titles": top,
"clients": [{"name": name, "minutes": round(value / 60, 1)} for name, value in sorted(clients.items(), key=lambda pair: -pair[1])[:6]],
"methods": [{"name": name, "minutes": round(value / 60, 1)} for name, value in sorted(methods.items(), key=lambda pair: -pair[1])],
"transcoding": {**{name: round(value, 1) for name, value in transcoding.items()},
"hardware": [{"name": name, "minutes": round(value, 1)} for name, value in sorted(hardware.items(), key=lambda pair: -pair[1])],
"audio_codecs": [{"name": name, "minutes": round(value, 1)} for name, value in sorted(audio_codecs.items(), key=lambda pair: -pair[1])],
"gpu_busy_minutes": None},
"recent": sorted(recent, key=lambda row: row["played_at"], reverse=True)[:20]}
async def get_insights(user: dict, days: int) -> dict:
runtime = await asyncio.to_thread(get_runtime_settings)
end = datetime.now(timezone.utc)
start = end - timedelta(days=days)
requests = await asyncio.to_thread(request_summary, user, start, end)
base = {"source": "Jellystat", "days": days, "timezone": "UTC", "requests": requests,
"is_admin": user.get("role") == "admin", "summary": None}
client = JellystatClient(runtime.jellystat_base_url, runtime.jellystat_api_key)
if not client.configured():
return {**base, "state": "not_configured"}
identity = await resolve_identity(user, runtime)
if not identity:
return {**base, "state": "unlinked"}
key = (runtime.jellystat_base_url, hashlib.sha256(runtime.jellystat_api_key.encode()).hexdigest(),
runtime.jellyfin_base_url, identity, days)
cached = _cache.get(key)
if cached and cached[0] > time.monotonic():
return {**base, **with_artwork(cached[1], user, runtime)}
history, libraries = await client.get_user_history(identity, start, end)
data = {**summarize(history, libraries, start, end), "state": "ready", "updated_at": end.isoformat(),
"period_start": start.isoformat(), "period_end": end.isoformat()}
for expired in [key for key, value in _cache.items() if value[0] <= time.monotonic()]:
_cache.pop(expired, None)
if len(_cache) >= 128:
_cache.pop(next(iter(_cache)))
_cache[key] = (time.monotonic() + CACHE_SECONDS, data)
return {**base, **with_artwork(data, user, runtime)}