"""Découverte : nouveautés des sources, incontournables et recommandations. Trois sections : - **latest** — « récemment ajoutés » scrapés sur chaque source activée (cliquables directement vers la fiche) ; - **must_watch** — titres les plus populaires du catalogue Kitsu (tous temps) ; - **for_you** — recommandations par genres : les genres des titres téléchargés (serveur), des favoris (par utilisateur) et des séries téléchargées sur Sonarr sont agrégés, puis Kitsu est interrogé sur ces catégories en excluant le déjà-possédé. Toutes les sources externes sont optionnelles : un échec (réseau, scraping, API) laisse la section vide et n'est jamais remonté au caller (dégradation gracieuse). Un cache mémoire TTL évite de re-scaper à chaque chargement de page. """ import asyncio import dataclasses import json import logging import re import time import unicodedata import httpx from app.config import get_settings from app.db import db from app.scrapers.base import ( ScrapeError, SourceScraper, all_sources, import_all_scrapers, ) from app.services.kitsu import KitsuService, normalize_title from app.services.settings import is_source_enabled from app.services.sonarr import sonarr logger = logging.getLogger(__name__) import_all_scrapers() # Bornes de l'algorithme _MAX_HISTORY_TITLES = 12 # titres récents analysés (téléchargements + favoris) _MAX_GENRES = 4 # genres retenus pour la requête Kitsu _KITSU_PAGE_MAX = 20 # limite dure de l'API Kitsu (page[limit] > 20 → 400) _ENRICH_CONCURRENCY = 6 # enrichissements Kitsu parallèles max (nouveautés) _LATEST_TTL_SECONDS = 600 # nouveautés : re-scrape au bout de 10 min _MUST_WATCH_TTL_SECONDS = 21600 # incontournables : quasi statique, 6 h _FOR_YOU_TTL_SECONDS = 3600 # recommandations : 1 h (l'historique évolue lentement) class _TTLCache: """Cache mémoire minimal avec expiration (mono-processus, suffisant ici).""" def __init__(self) -> None: self._data: dict[str, tuple[float, object]] = {} def get(self, key: str) -> object | None: entry = self._data.get(key) if entry is None: return None expires_at, value = entry if expires_at <= time.monotonic(): del self._data[key] return None return value def set(self, key: str, value: object, ttl_seconds: float) -> None: self._data[key] = (time.monotonic() + ttl_seconds, value) def clear(self, prefix: str = "") -> None: """Invalide les clés commençant par prefix (vide = tout le cache).""" for key in [k for k in self._data if k.startswith(prefix)]: del self._data[key] def category_slug(name: str) -> str: """Nom de genre → slug de catégorie Kitsu (« Slice of Life » → « slice-of-life »).""" decomposed = unicodedata.normalize("NFKD", name) ascii_only = "".join(char for char in decomposed if not unicodedata.combining(char)) return re.sub(r"[^a-z0-9]+", "-", ascii_only.casefold()).strip("-") class DiscoverService: """Agrégation des trois sections de découverte, avec cache mémoire.""" def __init__(self) -> None: self._cache = _TTLCache() self._kitsu = KitsuService() # ------------------------------------------------------------ nouveautés async def latest(self, limit: int = 24, allowed: set[str] | None = None) -> list[dict]: """Nouveautés toutes sources confondues, triées par date de sortie réelle. Les « récemment ajoutés » de chaque source sont fusionnés (doublons retirés), enrichis via Kitsu (date de début, statut de diffusion) puis triés du plus récent au plus ancien — ce qui sort / vient de sortir en premier. ``allowed`` restreint aux types de médias autorisés (préférence du compte), AVANT fusion/tri/troncature : sans lui, les titres hors Kitsu (séries, films réel) — sans date de sortie — seraient évincés du rail par la troncature. Le remplissage final équilibre les types présents (round-robin) pour qu'aucun ne soit écrasé par les autres en mode « les deux ». """ cache_key = f"latest:{limit}:{','.join(sorted(allowed)) if allowed else 'all'}" cached = self._cache.get(cache_key) if cached is not None: return cached # type: ignore[return-value] sources = [s for s in all_sources() if await is_source_enabled(s.name)] outcomes = await asyncio.gather(*(self._latest_of(source) for source in sources)) items = [item for outcome in outcomes for item in (outcome or [])] if allowed is not None: items = [item for item in items if item.get("media_type", "anime") in allowed] merged: dict[str, dict] = {} for item in items: key = normalize_title(item["title"]).casefold() existing = merged.get(key) if existing is None or (not existing.get("image_url") and item.get("image_url")): merged[key] = item semaphore = asyncio.Semaphore(_ENRICH_CONCURRENCY) async def bounded(item: dict) -> dict: async with semaphore: return await self._with_release_info(item) enriched = await asyncio.gather(*(bounded(item) for item in merged.values())) by_type: dict[str, list[dict]] = {} for item in sorted(enriched, key=lambda it: it.get("start_date") or "", reverse=True): by_type.setdefault(item.get("media_type", "anime"), []).append(item) result: list[dict] = [] pools = [list(pool) for pool in by_type.values()] while len(result) < limit and pools: for pool in pools[:]: result.append(pool.pop(0)) if len(result) >= limit: break if not pool: pools.remove(pool) self._cache.set(cache_key, result, _LATEST_TTL_SECONDS) return result async def _latest_of(self, source: SourceScraper) -> list[dict] | None: """Items latest() d'une source, aplatis avec les infos de source ([] si KO).""" try: results = await source.latest() except ScrapeError as exc: logger.warning("Nouveautés indisponibles pour %s : %s", source.name, exc) return None return [ {**dataclasses.asdict(r), "source": source.name, "label": source.label} for r in results ] async def _with_release_info(self, item: dict) -> dict: """Complète un item de nouveauté avec sa date de sortie Kitsu (None si absent).""" item.setdefault("start_date", None) item.setdefault("status", None) item.setdefault("rating", None) match = await self._kitsu_match_for_title(item["title"]) if match is None: return item attrs = match.get("attributes", {}) item["start_date"] = attrs.get("startDate") item["status"] = attrs.get("status") item["rating"] = KitsuService._to_rating_10(attrs.get("averageRating")) return item # --------------------------------------------------------- incontournables async def must_watch(self, limit: int = _KITSU_PAGE_MAX) -> list[dict]: """Titres les plus populaires du catalogue Kitsu (tous temps).""" limit = min(limit, _KITSU_PAGE_MAX) key = f"must_watch:{limit}" cached = self._cache.get(key) if cached is not None: return cached # type: ignore[return-value] items = await self._kitsu_anime({"sort": "-userCount", "page[limit]": limit}) self._cache.set(key, items, _MUST_WATCH_TTL_SECONDS) return items # ------------------------------------------------------------ pour toi async def for_you(self, user_id: int, limit: int = _KITSU_PAGE_MAX) -> dict: """Recommandations par genres, à partir de l'historique de l'utilisateur. Genres = téléchargements du serveur (titres → Kitsu) + favoris de l'utilisateur (genres du payload) + séries téléchargées sur Sonarr (genres fournis par Sonarr). On exclut les titres déjà possédés (local et Sonarr). """ limit = min(limit, _KITSU_PAGE_MAX) cache_key = f"for_you:{user_id}:{limit}" cached = self._cache.get(cache_key) if cached is not None: return cached # type: ignore[return-value] owned, favorite_genres = await self._owned(user_id) genre_counts = await self._genres_from_downloads(owned) for genre, count in favorite_genres.items(): genre_counts[genre] = genre_counts.get(genre, 0) + count sonarr_owned, sonarr_genres = await sonarr.profile() for genre, count in sonarr_genres.items(): genre_counts[genre] = genre_counts.get(genre, 0) + count owned |= sonarr_owned if not genre_counts: result: dict = {"based_on": [], "items": []} self._cache.set(cache_key, result, _FOR_YOU_TTL_SECONDS) return result top_genres = sorted(genre_counts, key=genre_counts.get, reverse=True)[:_MAX_GENRES] slugs = [category_slug(genre) for genre in top_genres] items = await self._kitsu_anime( { "filter[categories]": ",".join(slugs), "sort": "-userCount", "page[limit]": limit, # le déjà-possédé est filtré après } ) kept = [item for item in items if item["title"] and item["title"].casefold() not in owned] result = {"based_on": top_genres, "items": kept} self._cache.set(cache_key, result, _FOR_YOU_TTL_SECONDS) return result def invalidate_for_you(self) -> None: """Recommandations recalculées au prochain appel (réglages Sonarr modifiés).""" self._cache.clear("for_you:") async def _owned(self, user_id: int) -> tuple[set[str], dict[str, int]]: """Titres possédés (normalisés) + genres directement connus via les favoris.""" rows = await db.fetchall( "SELECT DISTINCT title FROM downloads ORDER BY created_at DESC LIMIT ?", (_MAX_HISTORY_TITLES,), ) fav_rows = await db.fetchall( "SELECT payload FROM favorites WHERE user_id = ? ORDER BY created_at DESC LIMIT ?", (user_id, _MAX_HISTORY_TITLES), ) owned = {normalize_title(row["title"]).casefold() for row in rows} owned.discard("") genre_counts: dict[str, int] = {} for row in fav_rows: try: payload = json.loads(row["payload"]) if row["payload"] else {} except (TypeError, ValueError): continue for genre in payload.get("genres") or []: if isinstance(genre, str) and genre.strip(): genre_counts[genre.strip()] = genre_counts.get(genre.strip(), 0) + 1 return owned, genre_counts async def _genres_from_downloads(self, titles: set[str]) -> dict[str, int]: """Genres Kitsu des titres téléchargés (cache DB puis recherche).""" semaphore = asyncio.Semaphore(_MAX_HISTORY_TITLES) async def genres_of(title: str) -> list[str]: async with semaphore: return await self._kitsu_genres_for_title(title) outcomes = await asyncio.gather(*(genres_of(t) for t in list(titles)[:_MAX_HISTORY_TITLES])) counts: dict[str, int] = {} for genres in outcomes: for genre in genres: counts[genre] = counts.get(genre, 0) + 1 return counts async def _kitsu_match_for_title(self, title: str) -> dict | None: """Match Kitsu d'un titre scrapé (cache DB 72 h via metadata_cache).""" query = normalize_title(title) if not query: return None cache_key = f"kitsu:anime:{query.casefold()}" match = await self._kitsu.get_cached(cache_key) if match is None: match = await self._kitsu.search_anime(title) if match is not None: await self._kitsu.set_cached(cache_key, match) return match async def _kitsu_genres_for_title(self, title: str) -> list[str]: """Genres Kitsu d'un titre scrapé (cache DB 72 h via metadata_cache). La recherche Kitsu ne renvoie plus les genres (`include=genres` vide) : on complète avec l'endpoint /anime//categories. """ match = await self._kitsu_match_for_title(title) if match is None: return [] genres = [g for g in match.get("genres", []) if isinstance(g, str)] if not genres: genres = await self._kitsu_categories(match.get("id")) match["genres"] = genres query = normalize_title(title) await self._kitsu.set_cached( # refresh avec les genres f"kitsu:anime:{query.casefold()}", match ) return genres # -------------------------------------------------------------- Kitsu async def _kitsu_anime(self, params: dict) -> list[dict]: """Requête générique liste Kitsu → items normalisés ([] si échec).""" settings = get_settings() try: async with httpx.AsyncClient( timeout=settings.http_timeout, headers={ "User-Agent": settings.user_agent, "Accept": "application/vnd.api+json", }, ) as client: response = await client.get(f"{settings.kitsu_base_url}/anime", params=params) response.raise_for_status() payload = response.json() except (httpx.HTTPError, ValueError) as exc: logger.warning("Liste Kitsu échouée (%s) : %s", params, exc) return [] return [ self._normalize_anime(item) for item in payload.get("data", []) if item.get("type") == "anime" ] async def _kitsu_categories(self, anime_id: object) -> list[str]: """Titres des catégories Kitsu d'un anime ([] si échec).""" if not anime_id: return [] settings = get_settings() try: async with httpx.AsyncClient( timeout=settings.http_timeout, headers={ "User-Agent": settings.user_agent, "Accept": "application/vnd.api+json", }, ) as client: response = await client.get( f"{settings.kitsu_base_url}/anime/{anime_id}/categories", params={"page[limit]": _KITSU_PAGE_MAX}, ) response.raise_for_status() payload = response.json() except (httpx.HTTPError, ValueError) as exc: logger.warning("Catégories Kitsu échouées (anime %s) : %s", anime_id, exc) return [] return [ attrs["title"] for item in payload.get("data", []) if isinstance(attrs := item.get("attributes", {}), dict) and attrs.get("title") ] @staticmethod def _normalize_anime(item: dict) -> dict: attrs = item.get("attributes", {}) titles = attrs.get("titles") or {} poster = attrs.get("posterImage") or {} return { "kitsu_id": item.get("id"), "title": attrs.get("canonicalTitle") or titles.get("en_jp"), "image_url": poster.get("large") or poster.get("medium") or poster.get("tiny"), "rating": KitsuService._to_rating_10(attrs.get("averageRating")), "year": KitsuService._extract_year(attrs.get("startDate")), "subtype": attrs.get("subtype"), "user_count": attrs.get("userCount"), } discover = DiscoverService()