diff --git a/.gitignore b/.gitignore index d6ad07c..edcd5b2 100644 --- a/.gitignore +++ b/.gitignore @@ -9,6 +9,8 @@ Cargo.lock # These are backup files generated by rustfmt **/*.rs.bk +__pycache__/ +*.pyc # MSVC Windows builds of rustc generate these, which store debugging information *.pdb diff --git a/docs/SSE.md b/docs/SSE.md index d021e16..3ffc3c2 100644 --- a/docs/SSE.md +++ b/docs/SSE.md @@ -231,24 +231,20 @@ interface MediasRatingMessage { rating: Rating; } -// Watched events (user-specific) -// IMPORTANT: The `id` field uses external IDs, NOT local database IDs. -// Format: "provider:id" (e.g., "imdb:tt1234567", "trakt:123456", "tmdb:550") -// For movies: Uses the best available external ID (priority: imdb > trakt > tmdb > tvdb > slug) -// For episodes: Uses external IDs or falls back to local "redseat:{id}" if no external IDs exist +// Watched events (user-specific). IDs include the media type and identity scheme. interface Watched { type: string; // MediaType: "movie", "episode", etc. - id: string; // External ID in format "provider:value" (e.g., "imdb:tt1234567") + id: string; // e.g. "movie:imdb/tt1234567" or "episode:redseat/seriesId/1/2" userRef?: string; date: number; // Timestamp when content was watched modified: number; } // Unwatched events (user-specific) -// NOTE: Different structure from Watched - contains ALL possible IDs for client matching +// NOTE: Different structure from Watched because the delete API accepts multiple IDs. interface Unwatched { type: string; // MediaType: "movie", "episode", etc. - ids: string[]; // All possible IDs in format "provider:value" (e.g., ["imdb:tt1234567", "trakt:12345", "tmdb:550"]) + ids: string[]; // History IDs marked as deleted userRef?: string; modified: number; } @@ -363,7 +359,7 @@ eventSource.addEventListener('unwatched', (event) => { const data: SseEvent = JSON.parse(event.data); if ('Unwatched' in data) { const unwatched = data.Unwatched; - // Unwatched events contain ALL possible IDs for the content + // Unwatched events contain the history IDs marked as deleted console.log(`Unmarked as watched: ${unwatched.type} with IDs: ${unwatched.ids.join(', ')}`); } }); @@ -573,14 +569,14 @@ Search endpoints support SSE streaming so clients receive results progressively Each SSE event has event type `results`. The data is a JSON object with a single key: the provider name, and the value is an array of results from that provider. -Results arrive one provider at a time. For series and movies, Trakt results are sent first, followed by each plugin (e.g., Anilist). For books, only plugin results are sent (no Trakt). +Results arrive one metadata plugin at a time. ### Event Format Each `results` event contains one provider's results: ```json -{"trakt": [{"metadata": {"serie": { ... }}, "images": []}]} +{"TMDB": [{"metadata": {"serie": { ... }}, "images": []}]} ``` Then a second event for the next provider: @@ -604,7 +600,7 @@ const resultsByProvider: Record = {}; eventSource.addEventListener('results', (event) => { const data = JSON.parse(event.data); - // data is e.g. { "trakt": [...] } or { "Anilist": [...] } + // data is e.g. { "TMDB": [...] } or { "Anilist": [...] } for (const [provider, results] of Object.entries(data)) { resultsByProvider[provider] = results; } @@ -633,7 +629,7 @@ Response format: ```json { - "trakt": [{"metadata": {"movie": { ... }}, "images": []}], + "TMDB": [{"metadata": {"movie": { ... }}, "images": []}], "Anilist": [{"metadata": {"movie": { ... }}, "images": [...]}] } ``` @@ -650,22 +646,16 @@ When a client falls behind and misses events (lag), the server will skip the mis ### Understanding the ID Format -The `watched` and `unwatched` events use **external IDs** (from providers like IMDb, Trakt, TMDb) rather than local database IDs. This allows watch history to be portable across different servers and sync with external services. +History IDs include the media type so IDs from different domains cannot collide. -**ID Format**: `provider:value` +| Content | Format | Example | +|---------|--------|---------| +| Movie with IMDb ID | `movie:imdb/` | `movie:imdb/tt1234567` | +| Movie without IMDb ID | `movie:redseat/` | `movie:redseat/abc123` | +| Series progress parent | `series:redseat/` | `series:redseat/series123` | +| Episode | `episode:redseat///` | `episode:redseat/series123/1/2` | -| Provider | Example | Content Types | -|----------|---------|---------------| -| `imdb` | `imdb:tt1234567` | Movies, Episodes | -| `trakt` | `trakt:123456` | Movies, Episodes, Series | -| `tmdb` | `tmdb:550` | Movies, Episodes, Series | -| `tvdb` | `tvdb:78901` | Episodes, Series | -| `slug` | `slug:the-matrix` | Movies, Series | -| `redseat` | `redseat:abc123` | Local fallback (episodes only) | - -**ID Selection Priority**: -- **Movies**: Uses the best external ID (priority: imdb > trakt > tmdb > slug) -- **Episodes**: Uses external IDs, or falls back to local `redseat:` ID if no external IDs exist +Episode IDs deliberately use the immutable local series ID and numeric season/episode tuple. Plugin metadata refreshes therefore cannot change watched state or progress IDs. ### REST API Endpoints @@ -681,30 +671,32 @@ The `watched` and `unwatched` events use **external IDs** (from providers like I { "date": 1705766400000 } ``` -**Direct History** (requires knowing the external ID): `POST /users/me/history` +**Direct History**: `POST /users/me/history` ```json { "type": "movie", - "id": "imdb:tt1234567", + "id": "movie:imdb/tt1234567", "date": 1705766400000 } ``` +For movie compatibility, `imdb:tt1234567` is also accepted and normalized to the typed form. Other media should send the typed history ID returned by the history API or SSE event. + #### Unmark as Watched (Remove from History) **Movies**: `DELETE /libraries/:libraryId/movies/:id/watched` **Episodes**: `DELETE /libraries/:libraryId/series/:serieId/seasons/:season/episodes/:number/watched` -**Direct History** (with multiple possible IDs): `DELETE /users/me/history` +**Direct History**: `DELETE /users/me/history` ```json { "type": "movie", - "ids": ["imdb:tt1234567", "trakt:12345", "tmdb:550"] + "ids": ["movie:imdb/tt1234567"] } ``` -The delete endpoints accept multiple IDs because the watched entry could have been created with any of the available external IDs. The server will try to delete entries matching any of the provided IDs. +The delete endpoint accepts an array so clients can delete more than one known history ID in one request. ### Example: Handling Watch State Changes @@ -736,25 +728,19 @@ eventSource.addEventListener('unwatched', (event) => { ### Matching SSE Events to Local Content -Since SSE events use external IDs, you need to match them against your local content's external IDs: +Match SSE events against the same typed history ID used by the REST API: ```typescript interface LocalMovie { id: string; // Local database ID imdb?: string; // "tt1234567" - trakt?: number; // 12345 - tmdb?: number; // 550 } -// For Watched events (single ID) function isMatchingWatchedEvent(movie: LocalMovie, eventId: string): boolean { - const [provider, value] = eventId.split(':'); - switch (provider) { - case 'imdb': return movie.imdb === value; - case 'trakt': return movie.trakt?.toString() === value; - case 'tmdb': return movie.tmdb?.toString() === value; - default: return false; - } + const historyId = movie.imdb + ? `movie:imdb/${movie.imdb}` + : `movie:redseat/${movie.id}`; + return eventId === historyId; } // For Unwatched events (array of IDs) @@ -823,14 +809,14 @@ syncHistory(); [ { "type": "movie", - "id": "imdb:tt1234567", + "id": "movie:imdb/tt1234567", "userRef": "user123", "date": 1705766400000, "modified": 1705852800000 }, { "type": "movie", - "id": "trakt:98765", + "id": "movie:imdb/tt9876543", "userRef": "user123", "date": 0, "modified": 1705939200000 diff --git a/scripts/remap_trakt_watched.py b/scripts/remap_trakt_watched.py new file mode 100644 index 0000000..e5fff6a --- /dev/null +++ b/scripts/remap_trakt_watched.py @@ -0,0 +1,546 @@ +#!/usr/bin/env python3 +from __future__ import annotations + +import argparse +import csv +import json +import shutil +import sqlite3 +import sys +import time +import urllib.error +import urllib.parse +import urllib.request +from dataclasses import dataclass +from pathlib import Path + + +USER_AGENT = ( + "Mozilla/5.0 (Windows NT 10.0; Win64; x64) " + "AppleWebKit/537.36 (KHTML, like Gecko) " + "Chrome/127.0.0.0 Safari/537.36" +) + + +@dataclass +class SeriesRow: + id: str + name: str + imdb: str | None + tmdb: int | None + trakt: int | None + tvdb: int | None + + +@dataclass +class EpisodeRow: + serie_ref: str + season: int + number: int + imdb: str | None + tmdb: int | None + trakt: int | None + tvdb: int | None + + @property + def redseat_id(self) -> str: + return f"episode:redseat/{self.serie_ref}/{self.season}/{self.number}" + + +class TraktClient: + def __init__(self, client_id: str) -> None: + self.client_id = client_id + + def _headers(self) -> dict[str, str]: + return { + "trakt-api-version": "2", + "trakt-api-key": self.client_id, + "Accept": "application/json", + "Content-Type": "application/json", + "User-Agent": USER_AGENT, + } + + def get_json(self, path: str, retries: int = 4) -> object: + url = f"https://api.trakt.tv{path}" + last_error: Exception | None = None + for attempt in range(retries): + request = urllib.request.Request(url, headers=self._headers()) + try: + with urllib.request.urlopen(request, timeout=30) as response: + payload = response.read().decode("utf-8") + return json.loads(payload) + except urllib.error.HTTPError as exc: + body = exc.read().decode("utf-8", errors="replace") + last_error = RuntimeError(f"{exc.code} for {url}: {body[:400]}") + if exc.code == 429: + retry_after = int(exc.headers.get("Retry-After", "1")) + time.sleep(retry_after) + continue + if exc.code >= 500 and attempt + 1 < retries: + time.sleep(1 + attempt) + continue + raise last_error + except (urllib.error.URLError, TimeoutError, json.JSONDecodeError) as exc: + last_error = exc + if attempt + 1 < retries: + time.sleep(1 + attempt) + continue + raise + if last_error: + raise last_error + raise RuntimeError(f"Unable to fetch {url}") + + def lookup_show_trakt_id(self, series: SeriesRow) -> int | None: + if series.trakt: + return series.trakt + + lookups = [] + if series.imdb: + lookups.append(f"/search/imdb/{urllib.parse.quote(series.imdb)}?type=show") + if series.tmdb: + lookups.append(f"/search/tmdb/{series.tmdb}?type=show") + if series.tvdb: + lookups.append(f"/search/tvdb/{series.tvdb}?type=show") + + for path in lookups: + result = self.get_json(path) + if isinstance(result, list) and result: + show = result[0].get("show", {}) + trakt_id = show.get("ids", {}).get("trakt") + if trakt_id: + return int(trakt_id) + return None + + def get_show_seasons(self, trakt_show_id: int) -> list[dict]: + result = self.get_json(f"/shows/{trakt_show_id}/seasons?extended=episodes") + return result if isinstance(result, list) else [] + + def get_movie(self, trakt_movie_id: int) -> dict: + result = self.get_json(f"/movies/{trakt_movie_id}") + return result if isinstance(result, dict) else {} + + +def load_tv_library(tv_db: Path) -> tuple[list[SeriesRow], dict[tuple[str, int, int], EpisodeRow]]: + conn = sqlite3.connect(tv_db) + conn.row_factory = sqlite3.Row + try: + series_rows = [ + SeriesRow( + id=row["id"], + name=row["name"] or "", + imdb=row["imdb"], + tmdb=row["tmdb"], + trakt=row["trakt"], + tvdb=row["tvdb"], + ) + for row in conn.execute( + "SELECT id, name, imdb, tmdb, trakt, tvdb FROM series ORDER BY name" + ) + ] + episodes = { + (row["serie_ref"], row["season"], row["number"]): EpisodeRow( + serie_ref=row["serie_ref"], + season=row["season"], + number=row["number"], + imdb=row["imdb"], + tmdb=row["tmdb"], + trakt=row["trakt"], + tvdb=row["tvdb"], + ) + for row in conn.execute( + """ + SELECT serie_ref, season, number, imdb, tmdb, trakt, tvdb + FROM episodes + """ + ) + } + return series_rows, episodes + finally: + conn.close() + + +def load_history(server_db: Path) -> tuple[list[sqlite3.Row], list[sqlite3.Row]]: + conn = sqlite3.connect(server_db) + conn.row_factory = sqlite3.Row + try: + watched = list( + conn.execute( + "SELECT type, id, user_ref, date, modified FROM Watched ORDER BY type, user_ref, id" + ) + ) + progress = list( + conn.execute( + "SELECT type, id, user_ref, progress, parent, modified FROM progress ORDER BY type, user_ref, id" + ) + ) + return watched, progress + finally: + conn.close() + + +def build_episode_map( + client: TraktClient, + series_rows: list[SeriesRow], + local_episodes: dict[tuple[str, int, int], EpisodeRow], + out_csv: Path, +) -> tuple[dict[str, str], list[str]]: + mapping: dict[str, str] = {} + unresolved_series: list[str] = [] + with out_csv.open("w", newline="", encoding="utf-8") as handle: + writer = csv.DictWriter( + handle, + fieldnames=[ + "source", + "trakt_episode_id", + "redseat_episode_id", + "serie_ref", + "serie_name", + "serie_imdb", + "serie_tmdb", + "serie_tvdb", + "serie_trakt", + "season", + "episode", + "episode_imdb", + "episode_tmdb", + "episode_tvdb", + ], + ) + writer.writeheader() + + for episode in local_episodes.values(): + if not episode.trakt: + continue + trakt_key = f"trakt:{episode.trakt}" + mapping[trakt_key] = episode.redseat_id + writer.writerow( + { + "source": "local_db", + "trakt_episode_id": trakt_key, + "redseat_episode_id": episode.redseat_id, + "serie_ref": episode.serie_ref, + "serie_name": "", + "serie_imdb": "", + "serie_tmdb": "", + "serie_tvdb": "", + "serie_trakt": "", + "season": episode.season, + "episode": episode.number, + "episode_imdb": episode.imdb or "", + "episode_tmdb": episode.tmdb or "", + "episode_tvdb": episode.tvdb or "", + } + ) + + for series in series_rows: + trakt_show_id = client.lookup_show_trakt_id(series) + if not trakt_show_id: + unresolved_series.append(f"{series.id}::{series.name}") + continue + + try: + seasons = client.get_show_seasons(trakt_show_id) + except Exception as exc: # noqa: BLE001 + unresolved_series.append(f"{series.id}::{series.name}::{exc}") + continue + + for season in seasons: + season_number = season.get("number") + if season_number is None: + continue + for api_episode in season.get("episodes", []): + episode_number = api_episode.get("number") + trakt_episode_id = api_episode.get("ids", {}).get("trakt") + if season_number is None or episode_number is None or not trakt_episode_id: + continue + + local_episode = local_episodes.get((series.id, season_number, episode_number)) + if not local_episode: + continue + + trakt_key = f"trakt:{trakt_episode_id}" + mapping[trakt_key] = local_episode.redseat_id + writer.writerow( + { + "source": "trakt_api", + "trakt_episode_id": trakt_key, + "redseat_episode_id": local_episode.redseat_id, + "serie_ref": series.id, + "serie_name": series.name, + "serie_imdb": series.imdb or "", + "serie_tmdb": series.tmdb or "", + "serie_tvdb": series.tvdb or "", + "serie_trakt": trakt_show_id, + "season": season_number, + "episode": episode_number, + "episode_imdb": api_episode.get("ids", {}).get("imdb") or "", + "episode_tmdb": api_episode.get("ids", {}).get("tmdb") or "", + "episode_tvdb": api_episode.get("ids", {}).get("tvdb") or "", + } + ) + return mapping, unresolved_series + + +def build_movie_map( + client: TraktClient, + history_rows: list[sqlite3.Row], + out_csv: Path, +) -> tuple[dict[str, str], list[str]]: + movie_ids = sorted( + { + int(row["id"].split(":", 1)[1]) + for row in history_rows + if row["type"] == "movie" and row["id"].startswith("trakt:") + } + ) + mapping: dict[str, str] = {} + unresolved: list[str] = [] + + with out_csv.open("w", newline="", encoding="utf-8") as handle: + writer = csv.DictWriter( + handle, + fieldnames=["trakt_movie_id", "imdb_movie_id", "title", "year"], + ) + writer.writeheader() + for trakt_movie_id in movie_ids: + try: + movie = client.get_movie(trakt_movie_id) + except Exception as exc: # noqa: BLE001 + unresolved.append(f"{trakt_movie_id}::{exc}") + continue + + imdb_id = movie.get("ids", {}).get("imdb") + if not imdb_id: + unresolved.append(str(trakt_movie_id)) + continue + + trakt_key = f"trakt:{trakt_movie_id}" + imdb_key = f"movie:imdb/{imdb_id}" + mapping[trakt_key] = imdb_key + writer.writerow( + { + "trakt_movie_id": trakt_key, + "imdb_movie_id": imdb_key, + "title": movie.get("title", ""), + "year": movie.get("year", ""), + } + ) + return mapping, unresolved + + +def rewrite_history( + server_db: Path, + episode_map: dict[str, str], + movie_map: dict[str, str], +) -> dict[str, int]: + backup_path = server_db.with_name( + f"{server_db.name}.bak-{time.strftime('%Y%m%d-%H%M%S')}" + ) + shutil.copy2(server_db, backup_path) + + conn = sqlite3.connect(server_db) + conn.row_factory = sqlite3.Row + stats = { + "backup_created": 1, + "episodes_updated": 0, + "episodes_merged": 0, + "episodes_unresolved": 0, + "movies_updated": 0, + "movies_merged": 0, + "movies_unresolved": 0, + "progress_updated": 0, + "progress_merged": 0, + "progress_unresolved": 0, + } + + try: + rows = list( + conn.execute( + "SELECT type, id, user_ref, date, modified FROM Watched ORDER BY type, user_ref, id" + ) + ) + conn.execute("BEGIN IMMEDIATE") + for row in rows: + kind = row["type"] + old_id = row["id"] + if kind == "episode": + new_id = episode_map.get(old_id) + unresolved_key = "episodes_unresolved" + updated_key = "episodes_updated" + merged_key = "episodes_merged" + elif kind == "movie": + new_id = movie_map.get(old_id) + unresolved_key = "movies_unresolved" + updated_key = "movies_updated" + merged_key = "movies_merged" + else: + continue + + if not new_id or new_id == old_id: + if old_id.startswith("trakt:"): + stats[unresolved_key] += 1 + continue + + destination = conn.execute( + """ + SELECT date, modified FROM Watched + WHERE type = ? AND id = ? AND user_ref = ? + """, + (kind, new_id, row["user_ref"]), + ).fetchone() + + if destination: + merged_date = ( + row["date"] + if row["modified"] >= destination["modified"] + else destination["date"] + ) + conn.execute( + """ + UPDATE Watched + SET date = ? + WHERE type = ? AND id = ? AND user_ref = ? + """, + (merged_date, kind, new_id, row["user_ref"]), + ) + stats[merged_key] += 1 + else: + conn.execute( + """ + INSERT INTO Watched (type, id, user_ref, date) + VALUES (?, ?, ?, ?) + """, + (kind, new_id, row["user_ref"], row["date"]), + ) + stats[updated_key] += 1 + conn.execute( + """ + UPDATE Watched SET date = 0 + WHERE type = ? AND id = ? AND user_ref = ? + """, + (kind, old_id, row["user_ref"]), + ) + + progress_rows = list( + conn.execute( + "SELECT type, id, user_ref, progress, parent, modified FROM progress ORDER BY type, user_ref, id" + ) + ) + for row in progress_rows: + kind = row["type"] + old_id = row["id"] + mapping = episode_map if kind == "episode" else movie_map if kind == "movie" else None + if mapping is None: + continue + new_id = mapping.get(old_id) + if not new_id or new_id == old_id: + if old_id.startswith("trakt:"): + stats["progress_unresolved"] += 1 + continue + + destination = conn.execute( + """ + SELECT progress, parent, modified FROM progress + WHERE type = ? AND id = ? AND user_ref = ? + """, + (kind, new_id, row["user_ref"]), + ).fetchone() + parent = row["parent"] + if kind == "episode": + parent = f"series:redseat/{new_id.split('/')[1]}" + if destination: + if destination["modified"] > row["modified"]: + value = destination["progress"] + parent = destination["parent"] or parent + else: + value = row["progress"] + conn.execute( + """ + UPDATE progress SET progress = ?, parent = ? + WHERE type = ? AND id = ? AND user_ref = ? + """, + (value, parent, kind, new_id, row["user_ref"]), + ) + stats["progress_merged"] += 1 + else: + conn.execute( + """ + INSERT INTO progress (type, id, user_ref, progress, parent) + VALUES (?, ?, ?, ?, ?) + """, + (kind, new_id, row["user_ref"], row["progress"], parent), + ) + stats["progress_updated"] += 1 + conn.execute( + "DELETE FROM progress WHERE type = ? AND id = ? AND user_ref = ?", + (kind, old_id, row["user_ref"]), + ) + conn.commit() + except Exception: + conn.rollback() + raise + finally: + conn.close() + + stats["backup_path"] = str(backup_path) # type: ignore[assignment] + return stats + + +def main() -> int: + parser = argparse.ArgumentParser( + description="Remap Redseat watched entries away from Trakt IDs." + ) + parser.add_argument("--server-db", required=True, type=Path) + parser.add_argument("--tv-db", required=True, type=Path) + parser.add_argument("--client-id", required=True) + parser.add_argument("--client-secret", default="") + parser.add_argument("--out-dir", type=Path, required=True) + args = parser.parse_args() + + args.out_dir.mkdir(parents=True, exist_ok=True) + if args.client_secret: + print("client_secret provided but not used by this script.", file=sys.stderr) + + watched_rows, progress_rows = load_history(args.server_db) + series_rows, local_episodes = load_tv_library(args.tv_db) + client = TraktClient(args.client_id) + + episode_csv = args.out_dir / "trakt_episode_map.csv" + movie_csv = args.out_dir / "trakt_movie_map.csv" + unresolved_series_txt = args.out_dir / "trakt_unresolved_series.txt" + unresolved_movies_txt = args.out_dir / "trakt_unresolved_movies.txt" + report_json = args.out_dir / "trakt_remap_report.json" + + episode_map, unresolved_series = build_episode_map( + client, series_rows, local_episodes, episode_csv + ) + movie_map, unresolved_movies = build_movie_map( + client, watched_rows + progress_rows, movie_csv + ) + stats = rewrite_history(args.server_db, episode_map, movie_map) + + unresolved_series_txt.write_text("\n".join(unresolved_series), encoding="utf-8") + unresolved_movies_txt.write_text("\n".join(unresolved_movies), encoding="utf-8") + + report = { + "series_count": len(series_rows), + "local_episode_count": len(local_episodes), + "episode_map_count": len(episode_map), + "movie_map_count": len(movie_map), + "unresolved_series_count": len(unresolved_series), + "unresolved_movie_count": len(unresolved_movies), + "history_rewrite": stats, + "files": { + "episode_csv": str(episode_csv), + "movie_csv": str(movie_csv), + "unresolved_series": str(unresolved_series_txt), + "unresolved_movies": str(unresolved_movies_txt), + }, + } + report_json.write_text(json.dumps(report, indent=2), encoding="utf-8") + print(json.dumps(report, indent=2)) + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/src/domain/watched.rs b/src/domain/watched.rs index e4e664b..5b9fe24 100644 --- a/src/domain/watched.rs +++ b/src/domain/watched.rs @@ -35,7 +35,7 @@ pub struct WatchedLight { pub struct WatchedForDelete { #[serde(rename = "type")] pub kind: MediaType, - /// Multiple possible IDs to try (imdb, trakt, tmdb, local, etc.) + /// History IDs to mark as deleted. pub ids: Vec, } @@ -44,7 +44,7 @@ pub struct WatchedForDelete { pub struct Unwatched { #[serde(rename = "type")] pub kind: MediaType, - /// All possible IDs for this content (imdb, trakt, tmdb, local, etc.) + /// History IDs marked as deleted. pub ids: Vec, pub user_ref: Option, pub modified: u64, diff --git a/src/model/episodes.rs b/src/model/episodes.rs index 4ead2c7..e40f7c8 100644 --- a/src/model/episodes.rs +++ b/src/model/episodes.rs @@ -1,10 +1,10 @@ -use std::collections::HashMap; +use std::collections::{HashMap, HashSet}; use async_recursion::async_recursion; use nanoid::nanoid; use query_external_ip::SourceError; use rs_plugin_common_interfaces::{ - domain::rs_ids::RsIds, + domain::rs_ids::{ApplyRsIds, RsIds}, lookup::{RsLookupEpisode, RsLookupMetadataResult, RsLookupQuery}, ImageType, MediaType, }; @@ -32,6 +32,7 @@ use crate::{ use super::{ entity_search::merge_result_ids, error::{Error, Result}, + history::{episode_history_id, episode_history_ids}, medias::{RsSort, RsSortOrder}, store::sql::SqlOrder, users::{ConnectedUser, HistoryQuery}, @@ -106,6 +107,22 @@ pub struct EpisodeForUpdate { } impl ModelController { + async fn episode_series_by_ref( + &self, + library_id: &str, + episodes: &[Episode], + ) -> RsResult> { + let store = self.store.get_library_store(library_id)?; + let serie_refs: HashSet = episodes.iter().map(|episode| episode.serie.clone()).collect(); + let mut series = HashMap::new(); + for serie_ref in serie_refs { + if let Some(serie) = store.get_serie(&serie_ref).await? { + series.insert(serie_ref, serie.item); + } + } + Ok(series) + } + async fn lookup_episodes_metadata( &self, library_id: &str, @@ -245,27 +262,35 @@ impl ModelController { requesting_user: &ConnectedUser, library_id: Option, ) -> RsResult<()> { - let ids: RsIds = episode.clone().into(); - let watched = self - .get_watched( - HistoryQuery { - types: vec![MediaType::Episode], - id: Some(ids.clone()), - ..Default::default() - }, - requesting_user, - library_id.clone(), - ) - .await?; - let progress = self - .get_view_progress(ids, requesting_user, library_id) - .await?; - if let Some(progress) = progress { - episode.progress = Some(progress.progress); - } - let watched = watched.first(); - if let Some(watched) = watched { - episode.watched = Some(watched.date); + if let Some(library_id) = library_id { + let store = self.store.get_library_store(&library_id)?; + if let Some(serie) = store.get_serie(&episode.serie).await? { + let history_ids = episode_history_ids(&serie.item, episode); + let watched = self + .get_watched( + HistoryQuery { + types: vec![MediaType::Episode], + id: Some(history_ids.clone()), + ..Default::default() + }, + requesting_user, + Some(library_id.clone()), + ) + .await?; + let progress = self + .get_view_progress( + history_ids, + requesting_user, + Some(library_id), + ) + .await?; + if let Some(progress) = progress { + episode.progress = Some(progress.progress); + } + if let Some(watched) = watched.first() { + episode.watched = Some(watched.date); + } + } } episode.fill_imdb_ratings(&self.imdb).await; Ok(()) @@ -276,6 +301,11 @@ impl ModelController { requesting_user: &ConnectedUser, library_id: Option, ) -> RsResult<()> { + let series_by_ref = if let Some(library_id) = library_id.as_deref() { + self.episode_series_by_ref(library_id, episodes).await? + } else { + HashMap::new() + }; let watched = self .get_watched( HistoryQuery { @@ -304,15 +334,12 @@ impl ModelController { .collect::>(); for episode in episodes { - let ids = RsIds::from(episode.clone()); - let ids_string: Vec = ids.into(); - for id in ids_string { - let watch = watched.get(&id); - if let Some(watch) = watch { + if let Some(serie) = series_by_ref.get(&episode.serie) { + let history_ids = episode_history_ids(serie, episode).as_all_ids(); + if let Some(watch) = history_ids.iter().find_map(|id| watched.get(id)) { episode.watched = Some(*watch); } - let progress = progresses.get(&id); - if let Some(progress) = progress { + if let Some(progress) = history_ids.iter().find_map(|id| progresses.get(id)) { episode.progress = Some(*progress); } } @@ -520,6 +547,20 @@ impl ModelController { let ids = self .get_serie_ids(library_id, serie_id, requesting_user) .await?; + let existing_episodes = self + .get_episodes( + library_id, + EpisodeQuery { + serie_ref: Some(serie_id.to_string()), + ..Default::default() + }, + requesting_user, + ) + .await?; + let existing_ids_by_episode: HashMap<(u32, u32), RsIds> = existing_episodes + .into_iter() + .map(|episode| ((episode.season, episode.number), RsIds::from(episode))) + .collect(); let mut all_episodes = self .lookup_episodes_metadata(library_id, serie_id, &ids, requesting_user) .await?; @@ -530,6 +571,15 @@ impl ModelController { ))) .into()); } + for episode in &mut all_episodes { + if let Some(existing_ids) = + existing_ids_by_episode.get(&(episode.season, episode.number)) + { + let mut merged_ids = RsIds::from(episode.clone()); + merged_ids.merge(existing_ids); + episode.apply_rs_ids(&merged_ids); + } + } let store = self.store.get_library_store_optional(library_id).ok_or( Error::LibraryStoreNotFoundFor( library_id.clone().to_string(), diff --git a/src/model/history.rs b/src/model/history.rs new file mode 100644 index 0000000..f047d29 --- /dev/null +++ b/src/model/history.rs @@ -0,0 +1,337 @@ +use std::collections::{HashMap, HashSet}; + +use rs_plugin_common_interfaces::{domain::rs_ids::RsIds, MediaType}; + +use crate::domain::{episode::Episode, library::LibraryType, movie::Movie, serie::Serie}; +use crate::tools::log::{log_info, LogServiceType}; + +use super::{ + episodes::EpisodeQuery, + movies::MovieQuery, + series::SerieQuery, + store::sql::users::{HistoryIdRewrite, ProgressIdRewrite}, + ModelController, +}; + +const HISTORY_MIGRATION: &str = "canonical_history_ids_v2"; + +pub fn series_history_id(serie: &Serie) -> String { + format!("series:redseat/{}", serie.id) +} + +pub fn movie_history_id(movie: &Movie) -> String { + movie + .imdb + .as_ref() + .map(|imdb| format!("movie:imdb/{imdb}")) + .unwrap_or_else(|| format!("movie:redseat/{}", movie.id)) +} + +pub fn episode_history_id(serie: &Serie, episode: &Episode) -> String { + format!( + "episode:redseat/{}/{}/{}", + serie.id, episode.season, episode.number + ) +} + +fn ids_with_history_id(mut ids: RsIds, history_id: String) -> RsIds { + ids.try_add(history_id) + .expect("history ids must stay in key:value format"); + ids +} + +pub fn movie_history_ids(movie: &Movie) -> RsIds { + ids_with_history_id(movie.clone().into(), movie_history_id(movie)) +} + +pub fn episode_history_ids(serie: &Serie, episode: &Episode) -> RsIds { + ids_with_history_id(episode.clone().into(), episode_history_id(serie, episode)) +} + +pub fn normalize_history_id(kind: &MediaType, id: String) -> String { + match kind { + MediaType::Movie if id.starts_with("movie:") => id, + MediaType::Movie => id + .strip_prefix("imdb:") + .map(|value| format!("movie:imdb/{value}")) + .or_else(|| { + id.strip_prefix("redseat:") + .map(|value| format!("movie:redseat/{value}")) + }) + .unwrap_or(id), + MediaType::Episode if id.starts_with("episode:") => id, + MediaType::Episode => id + .strip_prefix("redseat:") + .and_then(normalize_legacy_episode_id) + .unwrap_or(id), + _ => id, + } +} + +pub fn direct_history_ids(id: String) -> crate::Result { + let mut ids = RsIds::try_from(id.clone())?; + for kind in [MediaType::Movie, MediaType::Episode] { + let normalized = normalize_history_id(&kind, id.clone()); + if normalized != id { + ids.try_add(normalized)?; + } + } + Ok(ids) +} + +fn normalize_legacy_episode_id(value: &str) -> Option { + let mut parts = value.rsplitn(3, 'x'); + let number = parts.next()?.parse::().ok()?; + let season = parts.next()?.parse::().ok()?; + let serie = parts.next()?; + Some(format!("episode:redseat/{serie}/{season}/{number}")) +} + +fn typed_alias(kind: &str, id: &str) -> Option { + let (provider, value) = id.split_once(':')?; + Some(format!("{kind}:{provider}/{value}")) +} + +fn movie_legacy_ids(movie: &Movie) -> HashSet { + let ids: RsIds = movie.clone().into(); + ids.as_all_ids() + .into_iter() + .flat_map(|id| [Some(id.clone()), typed_alias("movie", &id)]) + .flatten() + .collect() +} + +fn episode_legacy_ids(serie: &Serie, episode: &Episode) -> HashSet { + let ids: RsIds = episode.clone().into(); + let mut aliases: HashSet = ids + .as_all_ids() + .into_iter() + .flat_map(|id| [Some(id.clone()), typed_alias("episode", &id)]) + .flatten() + .collect(); + + for key in [ + serie.imdb.as_ref().map(|id| format!("imdb/{id}")), + serie.tmdb.map(|id| format!("tmdb/{id}")), + serie.tvdb.map(|id| format!("tvdb/{id}")), + Some(format!("redseat/{}", serie.id)), + ] + .into_iter() + .flatten() + { + aliases.insert(format!( + "episode:{key}/{}/{}", + episode.season, episode.number + )); + } + aliases +} + +fn add_candidate( + candidates: &mut HashMap>, + old_id: String, + target: T, +) { + candidates.entry(old_id).or_default().insert(target); +} + +fn unique_candidates(candidates: HashMap>) -> HashMap { + candidates + .into_iter() + .filter_map(|(old_id, targets)| { + if targets.len() == 1 { + targets.into_iter().next().map(|target| (old_id, target)) + } else { + None + } + }) + .collect() +} + +fn is_current_history_id(kind: &MediaType, id: &str) -> bool { + match kind { + MediaType::Movie => id.starts_with("movie:imdb/") || id.starts_with("movie:redseat/"), + MediaType::Episode => id.starts_with("episode:redseat/"), + _ => true, + } +} + +impl ModelController { + pub async fn migrate_history_ids(&self) -> crate::Result<()> { + if self + .store + .is_data_migration_complete(HISTORY_MIGRATION) + .await? + { + return Ok(()); + } + + let watched = self.store.get_all_watched().await?; + let progress = self.store.get_all_view_progress_rows().await?; + let legacy_movie_ids: HashSet = watched + .iter() + .map(|row| (&row.kind, &row.id)) + .chain(progress.iter().map(|row| (&row.kind, &row.id))) + .filter(|(kind, id)| **kind == MediaType::Movie && !is_current_history_id(kind, id)) + .map(|(_, id)| id.clone()) + .collect(); + let legacy_episode_ids: HashSet = watched + .iter() + .map(|row| (&row.kind, &row.id)) + .chain(progress.iter().map(|row| (&row.kind, &row.id))) + .filter(|(kind, id)| **kind == MediaType::Episode && !is_current_history_id(kind, id)) + .map(|(_, id)| id.clone()) + .collect(); + + if legacy_movie_ids.is_empty() && legacy_episode_ids.is_empty() { + self.store + .complete_data_migration(HISTORY_MIGRATION) + .await?; + return Ok(()); + } + + let libraries = self.store.get_libraries().await?; + let mut movie_candidates = HashMap::>::new(); + let mut episode_candidates = HashMap::>::new(); + + for library in libraries { + if !matches!(library.kind, LibraryType::Movies | LibraryType::Shows) { + continue; + } + let store = self.store.get_library_store(&library.id)?; + + if library.kind == LibraryType::Movies && !legacy_movie_ids.is_empty() { + for movie in store.get_movies(MovieQuery::default()).await? { + let canonical_id = movie_history_id(&movie); + for legacy_id in movie_legacy_ids(&movie) { + if legacy_id != canonical_id && legacy_movie_ids.contains(&legacy_id) { + add_candidate(&mut movie_candidates, legacy_id, canonical_id.clone()); + } + } + } + } + + if library.kind == LibraryType::Shows && !legacy_episode_ids.is_empty() { + let series_by_ref: HashMap = store + .get_series(SerieQuery::default()) + .await? + .into_iter() + .map(|serie| (serie.item.id.clone(), serie.item)) + .collect(); + for episode in store.get_episodes(EpisodeQuery::default()).await? { + let Some(serie) = series_by_ref.get(&episode.serie) else { + continue; + }; + let target = ( + episode_history_id(serie, &episode), + series_history_id(serie), + ); + for legacy_id in episode_legacy_ids(serie, &episode) { + if legacy_episode_ids.contains(&legacy_id) { + add_candidate(&mut episode_candidates, legacy_id, target.clone()); + } + } + } + } + } + + let movie_targets = unique_candidates(movie_candidates); + let episode_targets = unique_candidates(episode_candidates); + let watched_rewrites = watched + .into_iter() + .filter_map(|watched| { + let user_ref = watched.user_ref?; + let new_id = match watched.kind { + MediaType::Movie => movie_targets.get(&watched.id)?.clone(), + MediaType::Episode => episode_targets.get(&watched.id)?.0.clone(), + _ => return None, + }; + Some(HistoryIdRewrite { + kind: watched.kind, + old_id: watched.id, + new_id, + user_ref, + }) + }) + .collect(); + let progress_rewrites = progress + .into_iter() + .filter_map(|progress| { + let (new_id, new_parent) = match progress.kind { + MediaType::Movie => (movie_targets.get(&progress.id)?.clone(), None), + MediaType::Episode => { + let (id, parent) = episode_targets.get(&progress.id)?; + (id.clone(), Some(parent.clone())) + } + _ => return None, + }; + Some(ProgressIdRewrite { + kind: progress.kind, + old_id: progress.id, + new_id, + new_parent, + user_ref: progress.user_ref, + }) + }) + .collect(); + + let (watched_count, progress_count) = self + .store + .apply_history_rewrites(watched_rewrites, progress_rewrites) + .await?; + self.store + .complete_data_migration(HISTORY_MIGRATION) + .await?; + log_info( + LogServiceType::Database, + format!( + "History migration complete: watched rewrites={}, progress rewrites={}", + watched_count, progress_count + ), + ); + Ok(()) + } + + pub async fn migrate_movie_history_id( + &self, + old_id: String, + new_id: String, + ) -> crate::Result<()> { + if old_id == new_id { + return Ok(()); + } + let watched_rewrites = self + .store + .get_all_watched() + .await? + .into_iter() + .filter(|row| row.kind == MediaType::Movie && row.id == old_id) + .filter_map(|row| { + Some(HistoryIdRewrite { + kind: row.kind, + old_id: row.id, + new_id: new_id.clone(), + user_ref: row.user_ref?, + }) + }) + .collect(); + let progress_rewrites = self + .store + .get_all_view_progress_rows() + .await? + .into_iter() + .filter(|row| row.kind == MediaType::Movie && row.id == old_id) + .map(|row| ProgressIdRewrite { + kind: row.kind, + old_id: row.id, + new_id: new_id.clone(), + new_parent: None, + user_ref: row.user_ref, + }) + .collect(); + self.store + .apply_history_rewrites(watched_rewrites, progress_rewrites) + .await?; + Ok(()) + } +} diff --git a/src/model/mod.rs b/src/model/mod.rs index 2994802..84e646d 100644 --- a/src/model/mod.rs +++ b/src/model/mod.rs @@ -13,6 +13,7 @@ pub mod deleted; pub mod entity_images; pub mod entity_search; pub mod episodes; +pub mod history; pub mod media_progresses; pub mod media_ratings; pub mod medias; @@ -188,6 +189,7 @@ impl ModelController { }); mc.cache_update_all_libraries().await?; + mc.migrate_history_ids().await?; let scheduler = &mc.scheduler; scheduler.start(mc.clone()).await?; diff --git a/src/model/movies.rs b/src/model/movies.rs index 7b89c62..eb01612 100644 --- a/src/model/movies.rs +++ b/src/model/movies.rs @@ -33,6 +33,7 @@ use super::{ entity_images::EntityImageConfig, entity_search::merge_result_ids, error::{Error, Result}, + history::{movie_history_id, movie_history_ids}, store::sql::SqlOrder, users::{ConnectedUser, HistoryQuery}, ModelController, @@ -245,11 +246,14 @@ impl ModelController { library_id: Option, ) -> RsResult<()> { movie.fill_imdb_ratings(&self.imdb).await; - - let ids: RsIds = movie.clone().into(); + let history_ids = movie_history_ids(movie); let progress = self - .get_view_progress(ids, requesting_user, library_id.clone()) + .get_view_progress( + history_ids.clone(), + requesting_user, + library_id.clone(), + ) .await?; if let Some(progress) = progress { movie.progress = Some(progress.progress); @@ -259,7 +263,7 @@ impl ModelController { .get_watched( HistoryQuery { types: vec![MediaType::Movie], - id: Some(movie.clone().into()), + id: Some(history_ids), ..Default::default() }, requesting_user, @@ -305,18 +309,12 @@ impl ModelController { .map(|e| (e.id, e.date)) .collect::>(); for movie in movies { - let ids = RsIds::from(movie.clone()); - let ids_string: Vec = ids.into(); - - for id in ids_string { - let watch = watched.get(&id); - if let Some(watch) = watch { - movie.watched = Some(*watch); - } - let progress = progresses.get(&id); - if let Some(progress) = progress { - movie.progress = Some(*progress); - } + let history_ids = movie_history_ids(movie).as_all_ids(); + if let Some(watch) = history_ids.iter().find_map(|id| watched.get(id)) { + movie.watched = Some(*watch); + } + if let Some(progress) = history_ids.iter().find_map(|id| progresses.get(id)) { + movie.progress = Some(*progress); } movie.fill_imdb_ratings(&self.imdb).await; @@ -375,8 +373,18 @@ impl ModelController { } if update.has_update() { let store = self.store.get_library_store(library_id)?; + let old_movie = + store + .get_movie(&movie_id) + .await? + .ok_or(SourcesError::UnableToFindMovie( + library_id.to_string(), + movie_id.clone(), + "update_movie".to_string(), + ))?; + let old_history_id = movie_history_id(&old_movie); store.update_movie(&movie_id, update).await?; - let person = + let movie = store .get_movie(&movie_id) .await? @@ -385,14 +393,16 @@ impl ModelController { movie_id.to_string(), "update_movie".to_string(), ))?; + self.migrate_movie_history_id(old_history_id, movie_history_id(&movie)) + .await?; self.send_movie(MoviesMessage { library: library_id.to_string(), movies: vec![MovieWithAction { action: ElementAction::Updated, - movie: person.clone(), + movie: movie.clone(), }], }); - Ok(person) + Ok(movie) } else { let movie = self .get_movie(library_id, movie_id, requesting_user) diff --git a/src/model/store/sql/users.rs b/src/model/store/sql/users.rs index 3435096..2298d64 100644 --- a/src/model/store/sql/users.rs +++ b/src/model/store/sql/users.rs @@ -31,6 +31,23 @@ use super::{ SqlOrder, SqlWhereType, }; +#[derive(Debug, Clone)] +pub struct HistoryIdRewrite { + pub kind: MediaType, + pub old_id: String, + pub new_id: String, + pub user_ref: String, +} + +#[derive(Debug, Clone)] +pub struct ProgressIdRewrite { + pub kind: MediaType, + pub old_id: String, + pub new_id: String, + pub new_parent: Option, + pub user_ref: String, +} + #[derive(Debug, Serialize, Deserialize, Clone, Default)] pub struct WatchedQuery { #[serde(rename = "type")] @@ -527,6 +544,22 @@ impl SqliteStore { Ok(row) } + pub async fn get_all_view_progress_rows(&self) -> Result> { + let row = self + .server_store + .call(move |conn| { + let mut query = conn.prepare( + "SELECT type, id, user_ref, progress, parent, modified FROM progress", + )?; + let rows = query.query_map(params![], Self::row_to_view_progress)?; + let progresses: Vec = + rows.collect::, rusqlite::Error>>()?; + Ok(progresses) + }) + .await?; + Ok(row) + } + pub async fn get_view_progess( &self, ids: RsIds, @@ -587,4 +620,189 @@ impl SqliteStore { .await?; Ok(()) } + + pub async fn apply_history_rewrites( + &self, + watched_rewrites: Vec, + progress_rewrites: Vec, + ) -> Result<(usize, usize)> { + let row = self + .server_store + .call(move |conn| { + let tx = conn.transaction()?; + let mut watched_count = 0usize; + let mut progress_count = 0usize; + + for rewrite in watched_rewrites { + if rewrite.old_id == rewrite.new_id { + continue; + } + + let source = tx + .query_row( + "SELECT date, modified FROM Watched WHERE type = ? AND id = ? AND user_ref = ?", + params![rewrite.kind, rewrite.old_id, rewrite.user_ref], + |row| Ok((row.get::<_, i64>(0)?, row.get::<_, u64>(1)?)), + ) + .optional()?; + + let Some((source_date, source_modified)) = source else { + continue; + }; + + let existing = tx + .query_row( + "SELECT date, modified FROM Watched WHERE type = ? AND id = ? AND user_ref = ?", + params![rewrite.kind, rewrite.new_id, rewrite.user_ref], + |row| Ok((row.get::<_, i64>(0)?, row.get::<_, u64>(1)?)), + ) + .optional()?; + + if let Some((existing_date, existing_modified)) = existing { + let merged_date = if source_modified >= existing_modified { + source_date + } else { + existing_date + }; + tx.execute( + "UPDATE Watched SET date = ? WHERE type = ? AND id = ? AND user_ref = ?", + params![ + merged_date, + rewrite.kind, + rewrite.new_id, + rewrite.user_ref + ], + )?; + } else { + tx.execute( + "INSERT INTO Watched (type, id, user_ref, date) VALUES (?, ?, ?, ?)", + params![rewrite.kind, rewrite.new_id, rewrite.user_ref, source_date], + )?; + } + tx.execute( + "UPDATE Watched SET date = 0 WHERE type = ? AND id = ? AND user_ref = ?", + params![rewrite.kind, rewrite.old_id, rewrite.user_ref], + )?; + watched_count += 1; + } + + for rewrite in progress_rewrites { + if rewrite.old_id == rewrite.new_id { + continue; + } + + let source = tx + .query_row( + "SELECT progress, parent, modified FROM progress WHERE type = ? AND id = ? AND user_ref = ?", + params![rewrite.kind, rewrite.old_id, rewrite.user_ref], + |row| { + Ok(( + row.get::<_, u64>(0)?, + row.get::<_, Option>(1)?, + row.get::<_, u64>(2)?, + )) + }, + ) + .optional()?; + + let Some((source_progress, source_parent, source_modified)) = source else { + continue; + }; + + let existing = tx + .query_row( + "SELECT progress, parent, modified FROM progress WHERE type = ? AND id = ? AND user_ref = ?", + params![rewrite.kind, rewrite.new_id, rewrite.user_ref], + |row| { + Ok(( + row.get::<_, u64>(0)?, + row.get::<_, Option>(1)?, + row.get::<_, u64>(2)?, + )) + }, + ) + .optional()?; + + if let Some((existing_progress, existing_parent, existing_modified)) = existing { + let (merged_progress, merged_parent) = + if source_modified >= existing_modified { + ( + source_progress, + rewrite + .new_parent + .clone() + .or(source_parent.clone()) + .or(existing_parent.clone()), + ) + } else { + ( + existing_progress, + rewrite.new_parent.clone().or(existing_parent.clone()), + ) + }; + tx.execute( + "UPDATE progress SET progress = ?, parent = ? WHERE type = ? AND id = ? AND user_ref = ?", + params![ + merged_progress, + merged_parent, + rewrite.kind, + rewrite.new_id, + rewrite.user_ref + ], + )?; + tx.execute( + "DELETE FROM progress WHERE type = ? AND id = ? AND user_ref = ?", + params![rewrite.kind, rewrite.old_id, rewrite.user_ref], + )?; + } else { + tx.execute( + "UPDATE progress SET id = ?, parent = ? WHERE type = ? AND id = ? AND user_ref = ?", + params![ + rewrite.new_id, + rewrite.new_parent, + rewrite.kind, + rewrite.old_id, + rewrite.user_ref + ], + )?; + } + progress_count += 1; + } + + tx.commit()?; + Ok((watched_count, progress_count)) + }) + .await?; + Ok(row) + } + + pub async fn is_data_migration_complete(&self, name: &str) -> Result { + let name = name.to_string(); + let complete = self + .server_store + .call(move |conn| { + let complete = conn.query_row( + "SELECT EXISTS(SELECT 1 FROM migrations WHERE name = ?)", + params![name], + |row| row.get(0), + )?; + Ok(complete) + }) + .await?; + Ok(complete) + } + + pub async fn complete_data_migration(&self, name: &str) -> Result<()> { + let name = name.to_string(); + self.server_store + .call(move |conn| { + conn.execute( + "INSERT INTO migrations (name, up, down) SELECT ?, '', '' WHERE NOT EXISTS (SELECT 1 FROM migrations WHERE name = ?)", + params![name, name], + )?; + Ok(()) + }) + .await?; + Ok(()) + } } diff --git a/src/model/users.rs b/src/model/users.rs index 26b8f14..0331dcf 100644 --- a/src/model/users.rs +++ b/src/model/users.rs @@ -24,6 +24,7 @@ use crate::{ use super::{ error::{Error, Result}, + history::{direct_history_ids, normalize_history_id}, libraries::ServerLibraryForRead, medias::RsSort, store::sql::{users::WatchedQuery, SqlOrder}, @@ -527,11 +528,12 @@ impl ModelController { pub async fn add_watched( &self, - watched: WatchedForAdd, + mut watched: WatchedForAdd, user: &ConnectedUser, library_id: Option, ) -> RsResult<()> { user.check_role(&UserRole::Read)?; + watched.id = normalize_history_id(&watched.kind, watched.id); let user_id = user.user_id()?; let modified = now().timestamp_millis() as u64; @@ -572,11 +574,16 @@ impl ModelController { /// Removes watched entries. Tries all provided IDs and emits a single SSE event with all IDs. pub async fn remove_watched( &self, - watched: WatchedForDelete, + mut watched: WatchedForDelete, user: &ConnectedUser, library_id: Option, ) -> RsResult<()> { user.check_role(&UserRole::Read)?; + watched.ids = watched + .ids + .into_iter() + .map(|id| normalize_history_id(&watched.kind, id)) + .collect(); let user_id = user.user_id()?; let modified = now().timestamp_millis() as u64; @@ -613,11 +620,12 @@ impl ModelController { pub async fn add_view_progress( &self, - progress: ViewProgressForAdd, + mut progress: ViewProgressForAdd, user: &ConnectedUser, library_id: Option, ) -> RsResult<()> { user.check_role(&UserRole::Read)?; + progress.id = normalize_history_id(&progress.kind, progress.id); let user_id = user.user_id()?; if let Some(library_id) = library_id { @@ -661,7 +669,7 @@ impl ModelController { id: String, user: &ConnectedUser, ) -> RsResult> { - let media_id = RsIds::try_from(id)?; + let media_id = direct_history_ids(id)?; self.get_view_progress(media_id, user, None).await } diff --git a/src/routes/episodes.rs b/src/routes/episodes.rs index 12547a2..a43bd88 100644 --- a/src/routes/episodes.rs +++ b/src/routes/episodes.rs @@ -14,6 +14,7 @@ use crate::{ error::RsError, model::{ episodes::{EpisodeForUpdate, EpisodeQuery}, + history::{episode_history_id, episode_history_ids, series_history_id}, medias::MediaQuery, users::{ConnectedUser, HistoryQuery}, ModelController, @@ -418,7 +419,19 @@ async fn handler_progress_get( .get_episode(&library_id, serie_id, season, number, &user) .await?; let progress = mc - .get_view_progress(episode.into(), &user, Some(library_id.to_string())) + .get_view_progress( + episode_history_ids( + &mc.get_serie(&library_id, episode.serie.clone(), &user) + .await? + .ok_or(Error::NotFound(format!( + "Unable to find serie for handler_progress_get" + )))? + .item, + &episode, + ), + &user, + Some(library_id.to_string()), + ) .await? .ok_or(Error::NotFound(format!( "Unable to get best view progress for handler_progress_get" @@ -443,21 +456,11 @@ async fn handler_progress_set( serie_id.to_string(), "handler_lookup_season".to_string(), ))?; - let id = RsIds::from(episode) - .into_best_external() - .ok_or(Error::NotFound(format!( - "Unable to get best external for handler_progress_set" - )))?; - let serie_id = RsIds::from(serie.item) - .into_best_external() - .ok_or(Error::NotFound(format!( - "Unable to get best external for handler_progress_set serie" - )))?; let progress = ViewProgressForAdd { kind: MediaType::Episode, - id, + id: episode_history_id(&serie.item, &episode), progress: progress.progress, - parent: Some(serie_id), + parent: Some(series_history_id(&serie.item)), }; mc.add_view_progress(progress, &user, Some(library_id)) .await?; @@ -504,8 +507,12 @@ async fn handler_watched_get( let episode = mc .get_episode(&library_id, serie_id, season, number, &user) .await?; + let serie = mc + .get_serie(&library_id, episode.serie.clone(), &user) + .await? + .ok_or(Error::NotFound(format!("Unable to find serie for watched get")))?; let query = HistoryQuery { - id: Some(episode.into()), + id: Some(episode_history_ids(&serie.item, &episode)), ..Default::default() }; let progress = mc @@ -526,14 +533,15 @@ async fn handler_watched_set( let episode = mc .get_episode(&library_id, serie_id.clone(), season, number, &user) .await?; - let id = RsIds::from(episode) - .into_best_external_or_local() + let serie = mc + .get_serie(&library_id, episode.serie.clone(), &user) + .await? .ok_or(Error::NotFound(format!( - "Unable to get best external for handler_watched_set" + "Unable to find serie for handler_watched_set" )))?; let watched = WatchedForAdd { kind: MediaType::Episode, - id, + id: episode_history_id(&serie.item, &episode), date: watched.date, }; mc.add_watched(watched, &user, Some(library_id)).await?; @@ -549,11 +557,15 @@ async fn handler_watched_delete( let episode = mc .get_episode(&library_id, serie_id.clone(), season, number, &user) .await?; - let rs_ids = RsIds::from(episode); - let ids = rs_ids.into_all_external_or_local(); + let serie = mc + .get_serie(&library_id, episode.serie.clone(), &user) + .await? + .ok_or(Error::NotFound(format!( + "Unable to find serie for handler_watched_delete" + )))?; let watched = WatchedForDelete { kind: MediaType::Episode, - ids, + ids: vec![episode_history_id(&serie.item, &episode)], }; mc.remove_watched(watched, &user, Some(library_id)).await?; diff --git a/src/routes/movies.rs b/src/routes/movies.rs index 335642b..c556965 100644 --- a/src/routes/movies.rs +++ b/src/routes/movies.rs @@ -11,6 +11,7 @@ use crate::{ error::RsError, model::{ episodes::EpisodeQuery, + history::{movie_history_id, movie_history_ids}, medias::MediaQuery, movies::{MovieQuery, RsMovieSort}, store::sql::SqlOrder, @@ -341,7 +342,11 @@ async fn handler_progress_get( ) -> Result> { let movie = mc.get_movie(&library_id, movie_id, &user).await?; let progress = mc - .get_view_progress(movie.into(), &user, Some(library_id.to_string())) + .get_view_progress( + movie_history_ids(&movie), + &user, + Some(library_id.to_string()), + ) .await? .ok_or(Error::NotFound("unable to get movie progress".to_string()))?; Ok(Json(json!(progress))) @@ -354,14 +359,9 @@ async fn handler_progress_set( Json(progress): Json, ) -> Result<()> { let movie = mc.get_movie(&library_id, movie_id, &user).await?; - let id = RsIds::from(movie) - .into_best_external() - .ok_or(Error::NotFound( - "unable to get best external id for movie progress".to_string(), - ))?; let progress = ViewProgressForAdd { kind: MediaType::Movie, - id, + id: movie_history_id(&movie), progress: progress.progress, parent: None, }; @@ -407,7 +407,7 @@ async fn handler_watched_get( ) -> Result> { let movie = mc.get_movie(&library_id, movie_id, &user).await?; let query = HistoryQuery { - id: Some(movie.into()), + id: Some(movie_history_ids(&movie)), ..Default::default() }; let progress = mc @@ -426,14 +426,9 @@ async fn handler_watched_set( Json(watched): Json, ) -> Result<()> { let movie = mc.get_movie(&library_id, movie_id, &user).await?; - let id = RsIds::from(movie) - .into_best_external() - .ok_or(Error::NotFound( - "Unable to get best external id for movie".to_string(), - ))?; let watched = WatchedForAdd { kind: MediaType::Movie, - id, + id: movie_history_id(&movie), date: watched.date, }; mc.add_watched(watched, &user, Some(library_id)).await?; @@ -447,11 +442,9 @@ async fn handler_watched_delete( user: ConnectedUser, ) -> Result<()> { let movie = mc.get_movie(&library_id, movie_id, &user).await?; - let rs_ids = RsIds::from(movie); - let ids = rs_ids.into_all_external(); let watched = WatchedForDelete { kind: MediaType::Movie, - ids, + ids: vec![movie_history_id(&movie)], }; mc.remove_watched(watched, &user, Some(library_id)).await?;