diff --git a/backend/src/api/endpoints/library_api.rs b/backend/src/api/endpoints/library_api.rs index e22fa40ba..5947739f7 100644 --- a/backend/src/api/endpoints/library_api.rs +++ b/backend/src/api/endpoints/library_api.rs @@ -1,6 +1,6 @@ use crate::api::model::{AppState, EventMessage}; use axum::response::IntoResponse; -use log::{debug, error, info}; +use log::{debug, error, info, warn}; use std::sync::Arc; use serde_json::json; use shared::model::{LibraryScanRequest, LibraryScanSummary, LibraryScanSummaryStatus, LibraryStatus}; @@ -13,6 +13,17 @@ async fn scan_library( ) -> axum::response::Response { debug!("Library scan requested (force_rescan: {})", request.force_rescan); + let Some(_update_guard) = app_state.update_guard.try_library() else { + warn!("Library update already in progress; update skipped."); + let response = LibraryScanSummary { + status: LibraryScanSummaryStatus::Error, + message: "Library update already in progress.".to_string(), + result: None, + }; + let _ = app_state.event_manager.send_event(EventMessage::LibraryScanProgress(response)); + return (axum::http::StatusCode::BAD_REQUEST, axum::Json(json!({"error": "Library update already in progress.".to_string()}))).into_response(); + }; + // Check if Library is enabled let lib_config = match app_state.app_config.config.load().library.as_ref() { Some(config) if config.enabled => config.clone(), @@ -27,8 +38,9 @@ async fn scan_library( } }; - let client = app_state.http_client.load_full().as_ref().clone(); tokio::spawn(async move { + let client = app_state.http_client.load_full().as_ref().clone(); + // Create processor and run scan let processor = LibraryProcessor::new(lib_config, client); diff --git a/backend/src/api/endpoints/v1_api_playlist.rs b/backend/src/api/endpoints/v1_api_playlist.rs index db8f19ba7..b2d79da51 100644 --- a/backend/src/api/endpoints/v1_api_playlist.rs +++ b/backend/src/api/endpoints/v1_api_playlist.rs @@ -67,7 +67,8 @@ async fn playlist_update( let valid_targets = Arc::new(valid_targets); tokio::spawn({ async move { - playlist::exec_processing(&http_client, app_config, valid_targets, Some(event_manager), Some(playlist_state)).await; + playlist::exec_processing(&http_client, app_config, valid_targets, Some(event_manager), + Some(playlist_state), Some(app_state.update_guard.clone())).await; } }); axum::http::StatusCode::ACCEPTED.into_response() diff --git a/backend/src/api/main_api.rs b/backend/src/api/main_api.rs index 14d6bf901..da1a6115a 100644 --- a/backend/src/api/main_api.rs +++ b/backend/src/api/main_api.rs @@ -11,7 +11,7 @@ use crate::api::endpoints::xmltv_api::xmltv_api_register; use crate::api::endpoints::xtream_api::xtream_api_register; use crate::api::hdhomerun_proprietary::spawn_proprietary_tasks; use crate::api::hdhomerun_ssdp::spawn_ssdp_discover_task; -use crate::api::model::{create_cache, create_http_client, ActiveProviderManager, ActiveUserManager, AppState, CancelTokens, ConnectionManager, DownloadQueue, EventManager, HdHomerunAppState, PlaylistStorageState, SharedStreamManager}; +use crate::api::model::{create_cache, create_http_client, ActiveProviderManager, ActiveUserManager, AppState, CancelTokens, ConnectionManager, DownloadQueue, EventManager, HdHomerunAppState, PlaylistStorageState, SharedStreamManager, UpdateGuard}; use crate::api::scheduler::exec_scheduler; use crate::api::serve::serve; use crate::model::{AppConfig, Config, Healthcheck, ProcessTargets, RateLimitConfig}; @@ -110,7 +110,8 @@ async fn create_shared_data( event_manager, cancel_tokens: Arc::new(ArcSwap::from_pointee(CancelTokens::default())), playlists: Arc::new(PlaylistStorageState::new()), - geoip + geoip, + update_guard: UpdateGuard::new(), } } @@ -125,12 +126,13 @@ fn exec_update_on_boot( config.update_on_boot }; if update_on_boot { - let app_state_clone = Arc::clone(&app_state.app_config); + let app_config_clone = Arc::clone(&app_state.app_config); let targets_clone = Arc::clone(targets); let playlist_state = Arc::clone(&app_state.playlists); let client = client.clone(); + let update_guard = Some(app_state.update_guard.clone()); tokio::spawn(async move { - playlist::exec_processing(&client, app_state_clone, targets_clone, None, Some(playlist_state)).await; + playlist::exec_processing(&client, app_config_clone, targets_clone, None, Some(playlist_state), update_guard).await; }); } } diff --git a/backend/src/api/model/app_state.rs b/backend/src/api/model/app_state.rs index d27b894f6..12c82fc95 100644 --- a/backend/src/api/model/app_state.rs +++ b/backend/src/api/model/app_state.rs @@ -16,9 +16,10 @@ use std::collections::HashMap; use std::sync::atomic::AtomicI8; use std::sync::Arc; use std::time::Duration; -use tokio::sync::Mutex; +use tokio::sync::{Mutex}; use tokio::task; use tokio_util::sync::CancellationToken; +use crate::api::model::update_guard::UpdateGuard; use crate::repository::storage::get_geoip_path; use crate::utils::GeoIp; @@ -263,6 +264,7 @@ pub struct AppState { pub cancel_tokens: Arc>, pub playlists: Arc, pub geoip: Arc>, + pub update_guard: UpdateGuard, } impl AppState { diff --git a/backend/src/api/model/mod.rs b/backend/src/api/model/mod.rs index def29e1a7..ac87198c8 100644 --- a/backend/src/api/model/mod.rs +++ b/backend/src/api/model/mod.rs @@ -13,6 +13,7 @@ mod event_manager; mod playlist_mem_cache; mod provider_lineup_manager; mod connection_manager; +mod update_guard; pub(in crate::api) use self::active_provider_manager::*; pub(in crate::api) use self::active_user_manager::*; @@ -29,3 +30,4 @@ pub use self::stream::*; pub(in crate::api) use self::stream_error::*; pub(crate) use self::streams::*; pub(in crate::api) use self::xtream::*; +pub use self::update_guard::*; diff --git a/backend/src/api/model/update_guard.rs b/backend/src/api/model/update_guard.rs new file mode 100644 index 000000000..48da1ef20 --- /dev/null +++ b/backend/src/api/model/update_guard.rs @@ -0,0 +1,43 @@ +use std::sync::Arc; +use tokio::sync::{OwnedSemaphorePermit, Semaphore}; + +#[derive(Debug, Clone)] +pub struct UpdateGuard { + playlist: Arc, + library: Arc, +} + +impl Default for UpdateGuard { + fn default() -> Self { + Self { + playlist: Arc::new(Semaphore::new(1)), + library: Arc::new(Semaphore::new(1)), + } + } +} + +impl UpdateGuard { + pub fn new() -> Self { + Self::default() + } + + pub fn try_playlist(&self) -> Option { + self.playlist + .clone() + .try_acquire_owned() + .ok() + .map(|permit| UpdateGuardPermit { _permit: permit }) + } + + pub fn try_library(&self) -> Option { + self.library + .clone() + .try_acquire_owned() + .ok() + .map(|permit| UpdateGuardPermit { _permit: permit }) + } +} + +pub struct UpdateGuardPermit { + _permit: OwnedSemaphorePermit, +} diff --git a/backend/src/api/scheduler.rs b/backend/src/api/scheduler.rs index bdf97f9bf..011234214 100644 --- a/backend/src/api/scheduler.rs +++ b/backend/src/api/scheduler.rs @@ -60,7 +60,8 @@ async fn start_scheduler(client: reqwest::Client, expression: &str, app_state: A let app_config = Arc::clone(&app_state.app_config); let event_manager = Arc::clone(&app_state.event_manager); let playlist_state = app_state.playlists.clone(); - exec_processing(&client, app_config, Arc::clone(&targets), Some(event_manager), Some(playlist_state)).await; + exec_processing(&client, app_config, Arc::clone(&targets), Some(event_manager), + Some(playlist_state), Some(app_state.update_guard.clone())).await; } () = cancel.cancelled() => { break; diff --git a/backend/src/library/nfo_reader.rs b/backend/src/library/nfo_reader.rs index c91ed1b58..f2e17c2fc 100644 --- a/backend/src/library/nfo_reader.rs +++ b/backend/src/library/nfo_reader.rs @@ -5,6 +5,16 @@ use quick_xml::Reader; use std::path::Path; use tokio::fs; +macro_rules! push_to_field_list { + ($field:expr, $current:expr) => { + if let Some(field) = $field.as_mut() { + field.push($current.clone()); + } else { + $field = Some(vec![$current.clone()]); + } + }; +} + /// NFO reader for parsing Kodi/Jellyfin/Emby/Plex metadata files pub struct NfoReader; @@ -116,37 +126,17 @@ impl NfoReader { "id" | "imdb" | "imdbid" => movie.imdb_id = Some(current_text.clone()), "tmdbid" => movie.tmdb_id = current_text.parse().ok(), "rating" => movie.rating = current_text.parse().ok(), - "genre" => if let Some(genres) = movie.genres.as_mut() { - genres.push(current_text.clone()); - } else { - movie.genres = Some(vec![current_text.clone()]); - }, - "director" => if let Some(field) = movie.directors.as_mut() { - field.push(current_text.clone()); - } else { - movie.directors = Some(vec![current_text.clone()]); - }, - "credits" | "writer" => if let Some(field) = movie.writers.as_mut() { - field.push(current_text.clone()); - } else { - movie.writers = Some(vec![current_text.clone()]); - }, - "studio" => if let Some(field) = movie.studios.as_mut() { - field.push(current_text.clone()); - } else { - movie.studios = Some(vec![current_text.clone()]); - }, + "genre" => push_to_field_list!(movie.genres, current_text), + "director" => push_to_field_list!(movie.directors, current_text), + "credits" | "writer" => push_to_field_list!(movie.writers, current_text), + "studio" => push_to_field_list!(movie.studios, current_text), "thumb" | "poster" => movie.poster = Some(current_text.clone()), "fanart" => movie.fanart = Some(current_text.clone()), "name" if in_actor => current_actor.name.clone_from(¤t_text), "role" if in_actor => current_actor.role = Some(current_text.clone()), "actor" => { if !current_actor.name.is_empty() { - if let Some(field) = movie.actors.as_mut() { - field.push(current_actor.clone()); - } else { - movie.actors = Some(vec![current_actor.clone()]); - } + push_to_field_list!(movie.actors, current_actor); } in_actor = false; } @@ -229,16 +219,8 @@ impl NfoReader { "tmdbid" => series.tmdb_id = current_text.parse().ok(), "tvdbid" => series.tvdb_id = current_text.parse().ok(), "rating" => series.rating = current_text.parse().ok(), - "genre" => if let Some(genres) = series.genres.as_mut() { - genres.push(current_text.clone()); - } else { - series.genres = Some(vec![current_text.clone()]); - }, - "studio" => if let Some(genres) = series.studios.as_mut() { - genres.push(current_text.clone()); - } else { - series.studios = Some(vec![current_text.clone()]); - }, + "genre" => push_to_field_list!(series.genres, current_text), + "studio" => push_to_field_list!(series.studios, current_text), "thumb" | "poster" => series.poster = Some(current_text.clone()), "fanart" => series.fanart = Some(current_text.clone()), "status" => series.status = Some(current_text.clone()), @@ -246,11 +228,7 @@ impl NfoReader { "role" if in_actor => current_actor.role = Some(current_text.clone()), "actor" => { if !current_actor.name.is_empty() { - if let Some(actors) = series.actors.as_mut() { - actors.push(current_actor.clone()); - } else { - series.actors = Some(vec![current_actor.clone()]); - } + push_to_field_list!(series.actors, current_actor); } in_actor = false; } diff --git a/backend/src/library/processor.rs b/backend/src/library/processor.rs index df792c19a..2fd012433 100644 --- a/backend/src/library/processor.rs +++ b/backend/src/library/processor.rs @@ -19,11 +19,11 @@ pub struct LibraryProcessor { impl LibraryProcessor { /// Creates a new Library processor from application config pub fn from_app_config(app_config: &AppConfig) -> Option { - let client = create_http_client(&app_config); + let client = create_http_client(app_config); app_config.config.load().library.as_ref().map(|lib_cfg| Self::new(lib_cfg.clone(), client)) } - /// Creates a new VOD processor with the given configuration + /// Creates a new Library processor with the given configuration pub fn new(config: LibraryConfig, client: reqwest::Client) -> Self { let storage_path = std::path::PathBuf::from(&config.metadata.path); let scanner = LibraryScanner::new(config.clone()); @@ -38,9 +38,9 @@ impl LibraryProcessor { } } - /// Performs a full VOD scan + /// Performs a full Library scan pub async fn scan(&self, force_rescan: bool) -> Result { - info!("Starting VOD scan (force_rescan: {force_rescan})"); + info!("Starting Library scan (force_rescan: {force_rescan})"); // Initialize storage self.storage.initialize().await?; diff --git a/backend/src/library/tmdb_client.rs b/backend/src/library/tmdb_client.rs index 07e65e57d..9246cce83 100644 --- a/backend/src/library/tmdb_client.rs +++ b/backend/src/library/tmdb_client.rs @@ -10,6 +10,23 @@ pub const TMDB_API_KEY: &str = "4219e299c89411838049ab0dab19ebd5"; const TMDB_API_BASE_URL: &str = "https://api.themoviedb.org/3"; const TMDB_IMAGE_BASE_URL: &str = "https://image.tmdb.org/t/p/w500"; + +// Helper function: Vec -> Option> +fn some_if_nonempty(v: Vec) -> Option> { + if v.is_empty() { None } else { Some(v) } +} + +// helper function: Crew-Filter +fn crew_names(credits: Option<&TmdbCredits>, jobs: &[&str]) -> Option> { + credits.as_ref().map(|c| { + c.crew + .iter() + .filter(|crew| jobs.contains(&crew.job.as_str())) + .map(|crew| crew.name.clone()) + .collect::>() + }).and_then(some_if_nonempty) +} + /// TMDB API client with rate limiting pub struct TmdbClient { api_key: String, @@ -83,122 +100,50 @@ impl TmdbClient { /// Fetches detailed movie information async fn fetch_movie_details(&self, movie_id: u32) -> Option { + sleep(Duration::from_millis(self.rate_limit_ms)).await; let url = format!("{TMDB_API_BASE_URL}/movie/{movie_id}?api_key={}&append_to_response=credits", self.api_key); - match self.client.get(&url).send().await { - Ok(response) => { - if response.status().is_success() { - match response.json::().await { - Ok(details) => Some(MediaMetadata::Movie(MovieMetadata { - title: details.title, - original_title: Some(details.original_title), - year: details.release_date.split('-').next().and_then(|y| y.parse().ok()), - plot: Some(details.overview), - tagline: details.tagline, - runtime: Some(details.runtime), - mpaa: None, // TMDB doesn't provide MPAA rating in basic response - imdb_id: details.imdb_id, - tmdb_id: Some(details.id), - tvdb_id: None, - rating: Some(details.vote_average), - genres: details.genres.as_ref().and_then(|list| { - let result: Vec = list.iter().map(|g| g.name.clone()).collect(); - if result.is_empty() { - None - } else { - Some(result) - } - }), - directors: details - .credits - .as_ref() - .and_then(|c| { - let list: Vec = c.crew - .iter() - .filter(|crew| crew.job == "Director") - .map(|crew| crew.name.clone()) - .collect(); - if list.is_empty() { - None - } else { - Some(list) - } - }), - writers: details - .credits - .as_ref() - .and_then(|c| { - let list: Vec = c.crew - .iter() - .filter(|crew| crew.job == "Writer" || crew.job == "Screenplay") - .map(|crew| crew.name.clone()) - .collect(); - if list.is_empty() { - None - } else { - Some(list) - } - }), - actors: details - .credits - .as_ref() - .and_then(|c| { - let actors: Vec = c.cast - .iter() - .take(10) // Limit to top 10 actors - .map(|actor| Actor { - name: actor.name.clone(), - role: Some(actor.character.clone()), - thumb: actor - .profile_path - .as_ref() - .map(|p| format!("{TMDB_IMAGE_BASE_URL}{p}")), - }) - .collect(); - - if actors.is_empty() { - None - } else { - Some(actors) - } - }), - studios: - details.production_companies.as_ref().and_then(|list| { - let result: Vec = list.iter().map(|n| n.name.clone()).collect(); - if result.is_empty() { - None - } else { - Some(result) - } - }), - poster: details - .poster_path - .map(|p| format!("{TMDB_IMAGE_BASE_URL}{p}")), - fanart: details - .backdrop_path - .map(|p| format!("{TMDB_IMAGE_BASE_URL}{p}")), - source: MetadataSource::Tmdb, - last_updated: chrono::Utc::now().timestamp(), - })), - Err(e) => { - error!("Failed to parse TMDB movie details: {e}"); - None - } - } - } else { - warn!("TMDB API error fetching movie details: {}", response.status()); - None - } - } - Err(e) => { - error!("TMDB API request failed: {e}"); - None - } + let response = self.client.get(&url).send().await.ok()?; + if !response.status().is_success() { + warn!("TMDB API error fetching movie details: {}", response.status()); + return None; } + + let details: TmdbMovieDetails = response.json().await.ok()?; + + Some(MediaMetadata::Movie(MovieMetadata { + title: details.title, + original_title: Some(details.original_title), + year: details.release_date.split('-').next().and_then(|y| y.parse().ok()), + plot: Some(details.overview), + tagline: details.tagline, + runtime: Some(details.runtime), + mpaa: None, + imdb_id: details.imdb_id, + tmdb_id: Some(details.id), + tvdb_id: None, + rating: Some(details.vote_average), + genres: details.genres.as_ref().map(|list| list.iter().map(|g| g.name.clone()).collect()).and_then(some_if_nonempty), + directors: crew_names(details.credits.as_ref(), &["Director"]), + writers: crew_names(details.credits.as_ref(), &["Writer", "Screenplay"]), + actors: details.credits.as_ref().and_then(|c| { + some_if_nonempty(c.cast.iter().take(10).map(|actor| Actor { + name: actor.name.clone(), + role: Some(actor.character.clone()), + thumb: actor.profile_path.as_ref().map(|p| format!("{TMDB_IMAGE_BASE_URL}{p}")), + }).collect::>()) + }), + studios: details.production_companies.as_ref().map(|list| list.iter().map(|n| n.name.clone()).collect()).and_then(some_if_nonempty), + poster: details.poster_path.map(|p| format!("{TMDB_IMAGE_BASE_URL}{p}")), + fanart: details.backdrop_path.map(|p| format!("{TMDB_IMAGE_BASE_URL}{p}")), + source: MetadataSource::Tmdb, + last_updated: chrono::Utc::now().timestamp(), + })) } + /// Searches for a TV series by title and optional year pub async fn search_series(&self, title: &str, year: Option) -> Option { sleep(Duration::from_millis(self.rate_limit_ms)).await; diff --git a/backend/src/main.rs b/backend/src/main.rs index f2e984c16..763e528d1 100644 --- a/backend/src/main.rs +++ b/backend/src/main.rs @@ -185,7 +185,7 @@ async fn start_in_cli_mode(cfg: Arc, targets: Arc) { error!("Failed to build client {err}"); reqwest::Client::new() }); - playlist::exec_processing(&client, cfg, targets, None, None).await; + playlist::exec_processing(&client, cfg, targets, None, None, None).await; } async fn start_in_server_mode(cfg: Arc, targets: Arc) { diff --git a/backend/src/processing/processor/playlist.rs b/backend/src/processing/processor/playlist.rs index 990d89f73..d824ace2d 100644 --- a/backend/src/processing/processor/playlist.rs +++ b/backend/src/processing/processor/playlist.rs @@ -6,10 +6,10 @@ use crate::Config; use std::collections::{HashMap, HashSet}; use std::path::PathBuf; use std::sync::Arc; -use tokio::task::JoinSet; use tokio::sync::Mutex; +use tokio::task::JoinSet; -use crate::api::model::{EventManager, EventMessage, PlaylistStorageState}; +use crate::api::model::{EventManager, EventMessage, PlaylistStorageState, UpdateGuard}; use crate::messaging::send_message_json; use crate::model::Epg; use crate::model::FetchedPlaylist; @@ -24,8 +24,8 @@ use crate::processing::processor::trakt::process_trakt_categories_for_target; use crate::processing::processor::xtream_series::playlist_resolve_series; use crate::processing::processor::xtream_vod::playlist_resolve_vod; use crate::repository::playlist_repository::persist_playlist; -use crate::utils::{debug_if_enabled, trace_if_enabled}; use crate::utils::StepMeasure; +use crate::utils::{debug_if_enabled, trace_if_enabled}; use deunicode::deunicode; use futures::StreamExt; use log::{debug, error, info, log_enabled, trace, warn, Level}; @@ -654,9 +654,25 @@ async fn process_watch(cfg: &Config, client: &reqwest::Client, target: &ConfigTa } } -pub async fn exec_processing(client: &reqwest::Client, app_config: Arc, targets: Arc, event_manager: Option>, playlist_state: Option>) { - let start_time = Instant::now(); +pub async fn exec_processing(client: &reqwest::Client, app_config: Arc, targets: Arc, + event_manager: Option>, playlist_state: Option>, + update_guard: Option) { + let _guard = if let Some(guard) = update_guard { + if let Some(permit) = guard.try_playlist() { + Some(permit) + } else { + warn!("Playlist update already in progress; update skipped."); + if let Some(events) = event_manager.as_ref() { + events.send_event(EventMessage::PlaylistUpdate(PlaylistUpdateState::Failure)); + } + return; + } + } else { + None + }; + let event_manager_clone = event_manager.clone(); + let start_time = Instant::now(); let (stats, errors) = process_sources(client, &app_config, targets.clone(), event_manager_clone, playlist_state.as_ref()).await; // log errors for err in &errors {