diff --git a/backend/src/api/endpoints/v1_api_playlist.rs b/backend/src/api/endpoints/v1_api_playlist.rs index 92e7adbae..f664de91d 100644 --- a/backend/src/api/endpoints/v1_api_playlist.rs +++ b/backend/src/api/endpoints/v1_api_playlist.rs @@ -77,7 +77,7 @@ async fn playlist_update( async move { exec_processing(&http_client, app_config, valid_targets, Some(event_manager), Some(playlist_state), Some(app_state.update_guard.clone()), - disabled_headers, Some(provider_manager), Some(metadata_manager)).await; + disabled_headers, Some(provider_manager), Some(metadata_manager), None, None).await; } }); axum::http::StatusCode::ACCEPTED.into_response() diff --git a/backend/src/api/endpoints/websocket_api.rs b/backend/src/api/endpoints/websocket_api.rs index 77d160533..592ea9d13 100644 --- a/backend/src/api/endpoints/websocket_api.rs +++ b/backend/src/api/endpoints/websocket_api.rs @@ -285,7 +285,7 @@ async fn handle_event_message(socket: &mut WebSocket, event: EventMessage, handl .await .map_err(|e| format!("Library scan progress event: {e} "))?; } - EventMessage::InputMetadataUpdatesCompleted(_) => { + EventMessage::InputMetadataUpdatesCompleted(_) | EventMessage::InputMetadataUpdatesStarted(_) => { // Internal event, ignore for websocket clients } } diff --git a/backend/src/api/main_api.rs b/backend/src/api/main_api.rs index 74c61ac6f..ab71ba0fb 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, create_http_client_no_redirect, ActiveProviderManager, ActiveUserManager, AppState, CancelTokens, ConnectionManager, DownloadQueue, EventManager, HdHomerunAppState, MetadataUpdateManager, PlaylistStorageState, SharedStreamManager, UpdateGuard}; +use crate::api::model::{create_cache, create_http_client, create_http_client_no_redirect, ActiveProviderManager, ActiveUserManager, AppState, CancelTokens, ConnectionManager, DownloadQueue, EventManager, EventMessage, HdHomerunAppState, MetadataUpdateManager, PlaylistStorageState, SharedStreamManager, UpdateGuard}; use crate::api::panel_api::sync_panel_api_exp_dates_on_boot; use crate::api::scheduler::{exec_interner_prune, exec_scheduler}; use crate::api::serve::serve; @@ -26,10 +26,11 @@ use arc_swap::{ArcSwap, ArcSwapOption}; use axum::extract::connect_info::ConnectInfo; use axum::Router; use axum::{extract::Request, middleware::Next}; -use log::{debug, error, info}; +use log::{debug, error, info, warn}; use shared::error::TuliproxError; use shared::utils::{concat_path_leading_slash, sanitize_sensitive_info}; use shared::{info_err, info_err_res}; +use std::collections::{HashMap, HashSet}; use std::net::SocketAddr; use std::path::PathBuf; use std::sync::atomic::AtomicI8; @@ -136,9 +137,10 @@ fn exec_update_on_boot( let disabled_headers = app_state.get_disabled_headers(); let provider_manager = Arc::clone(&app_state.active_provider); let metadata_manager = Arc::clone(&app_state.metadata_manager); + let event_manager = Some(Arc::clone(&app_state.event_manager)); tokio::spawn(async move { - exec_processing(&client, app_config_clone, targets_clone, None, Some(playlist_state), update_guard, disabled_headers, Some(provider_manager), Some(metadata_manager)).await; + exec_processing(&client, app_config_clone, targets_clone, event_manager, Some(playlist_state), update_guard, disabled_headers, Some(provider_manager), Some(metadata_manager), None, None).await; }); } } @@ -459,75 +461,124 @@ async fn log_req(req: Request, next: Next) -> impl axum::response::IntoResponse } -fn exec_input_update_listener(_app_state: &Arc, _targets: &Arc) { - // TODO this method starts an playlist update after the metadata is resolved. - // Currently to many playlist updates are started and therefore it is disabled. +fn exec_input_update_listener(app_state: &Arc, targets: &Arc) { + let app_state = Arc::clone(app_state); + let targets = Arc::clone(targets); - // let app_state = Arc::clone(app_state); - // let targets = Arc::clone(targets); + tokio::spawn(async move { + let mut rx = app_state.event_manager.get_event_channel(); + // Map>: Tracks which inputs are currently updating for a given target + let mut active_target_inputs: HashMap> = HashMap::new(); + // Set: Tracks which targets are pending an update once their inputs are done + let mut pending_targets: HashSet = HashSet::new(); - // tokio::spawn(async move { - // let mut rx = app_state.event_manager.get_event_channel(); - // let app_state = Arc::clone(&app_state); - // loop { - // match rx.recv().await { - // Ok(EventMessage::InputMetadataUpdatesCompleted(input_name)) => { - // // Find all targets that use this input - // let mut targets_to_update = HashSet::new(); - // let sources = app_state.app_config.sources.load(); - // - // for source in &sources.sources { - // if source.inputs.contains(&input_name) { - // for target in &source.targets { - // // Update strategy "Bundled" explicitly waits for this event. - // // Update strategy "Instant" triggers updates per item, but we ALSO trigger a full - // // update here to ensure everything is consistent (e.g. M3U files, cleaning up temp files). - // // This guarantees that any STRM/M3U files are generated/updated after metadata resolution. - // - // // Check if target matches process targets (CLI args or schedule) - // if let Some(valid_targets) = crate::api::scheduler::get_process_targets(&app_state.app_config, &targets, Some(&vec![target.name.clone()])).as_ref().into() { - // if valid_targets.enabled && !valid_targets.targets.is_empty() { - // targets_to_update.insert(target.name.clone()); - // } - // } - // } - // } - // } - // - // if !targets_to_update.is_empty() { - // let targets_to_update: Vec<_> = targets_to_update.into_iter().collect(); - // // Small delay to ensure any lingering updates or file locks from the background thread are fully released - // tokio::time::sleep(std::time::Duration::from_millis(500)).await; - // - // info!("Triggering playlist update for targets due to metadata change completion: {targets_to_update:?}"); - // - // let client = app_state.http_client.load().as_ref().clone(); - // let app_config = Arc::clone(&app_state.app_config); - // let event_manager = Arc::clone(&app_state.event_manager); - // let playlist_state = Arc::clone(&app_state.playlists); - // let disabled_headers = app_state.get_disabled_headers(); - // - // if let Ok(process_targets) = sources.validate_targets(Some(&targets_to_update)) { - // let proc_targets = Arc::new(process_targets); - // - // let update_guard = app_state.update_guard.clone(); - // - // let app_state_clone = app_state.clone(); - // tokio::spawn(async move { - // exec_processing( - // &client, app_config, proc_targets, Some(event_manager), - // Some(playlist_state), Some(update_guard), - // disabled_headers, - // Some(app_state_clone.active_provider.clone()), Some(app_state_clone.metadata_manager.clone()), - // ).await; - // }); - // } - // } - // } - // Ok(_) - // | Err(tokio::sync::broadcast::error::RecvError::Lagged(_)) => {}, - // Err(tokio::sync::broadcast::error::RecvError::Closed) => break, - // } - // } - // }); + loop { + match rx.recv().await { + Ok(EventMessage::InputMetadataUpdatesStarted(input_name)) => { + let sources = app_state.app_config.sources.load(); + for source in &sources.sources { + if source.inputs.iter().any(|i| i.as_ref() == input_name.as_ref()) { + for target in &source.targets { + // Check if this target is allowed by the global process targets + if targets.enabled && !targets.target_names.contains(&target.name) { + continue; + } + // Add this input to the active set for this target + active_target_inputs + .entry(target.name.clone()) + .or_default() + .insert(input_name.to_string()); + // Mark target as potentially needing an update + pending_targets.insert(target.name.clone()); + } + } + } + } + Ok(EventMessage::InputMetadataUpdatesCompleted(input_name)) => { + // 1. Remove this input from all active sets + let target_names: Vec = active_target_inputs.keys().cloned().collect(); + let mut targets_to_trigger = Vec::new(); + + for target_name in target_names { + if let std::collections::hash_map::Entry::Occupied(mut entry) = active_target_inputs.entry(target_name.clone()) { + let inputs = entry.get_mut(); + inputs.remove(input_name.as_ref()); + if inputs.is_empty() { + entry.remove(); + // If this target was pending, it's now ready to trigger + if pending_targets.remove(&target_name) { + targets_to_trigger.push(target_name); + } + } + } + } + + if !targets_to_trigger.is_empty() { + info!("Triggering playlist update for targets due to metadata change completion: {targets_to_trigger:?}"); + + let client = app_state.http_client.load().as_ref().clone(); + let app_config = Arc::clone(&app_state.app_config); + let event_manager = Arc::clone(&app_state.event_manager); + let playlist_state = Arc::clone(&app_state.playlists); + let disabled_headers = app_state.get_disabled_headers(); + let sources = app_config.sources.load(); + + // For each target, we need to gather ALL its inputs to pass as pre_processed_inputs + let targets_set: HashSet = targets_to_trigger.iter().map(Clone::clone).collect(); + + if let Ok(process_targets) = sources.validate_targets(Some(&targets_to_trigger)) { + let proc_targets = Arc::new(process_targets); + // Collect all inputs for these targets + let mut pre_processed_inputs: HashSet> = HashSet::new(); + for source in &sources.sources { + for target in &source.targets { + if targets_set.contains(&target.name) { + for input in &source.inputs { + pre_processed_inputs.insert(input.clone()); + } + } + } + } + + let update_guard = app_state.update_guard.clone(); + let app_state_clone = app_state.clone(); + + // SPAWN the trigger logic so we don't block the event loop with sleep or heavy processing setup + tokio::spawn(async move { + // Small delay to ensure any lingering updates or file locks from the background thread are fully released + tokio::time::sleep(std::time::Duration::from_millis(500)).await; + + // Wait for any current update to finish before starting a new one + let lock_opt = update_guard.acquire_playlist_lock().await; + + // Only proceed if we successfully acquired the lock (semaphore not closed) + if let Some(lock) = lock_opt { + exec_processing( + &client, app_config, proc_targets, Some(event_manager), + Some(playlist_state), Some(update_guard), + disabled_headers, + Some(app_state_clone.active_provider.clone()), Some(app_state_clone.metadata_manager.clone()), + Some(pre_processed_inputs), + Some(lock), // Pass the acquired permit to stay active during processing + ).await; + } else { + warn!("Skipping triggered update because shutdown signal received (lock closed)"); + } + }); + } else { + warn!("Failed to validate targets for triggered update: {targets_to_trigger:?}"); + } + } + } + Ok(_) => {} + Err(tokio::sync::broadcast::error::RecvError::Lagged(skipped)) => { + warn!("Input update listener lagged by {skipped} messages. Resetting tracking state (active_targets={}, pending={}) to avoid inconsistencies.", + active_target_inputs.len(), pending_targets.len()); + active_target_inputs.clear(); + pending_targets.clear(); + } + Err(tokio::sync::broadcast::error::RecvError::Closed) => break, + } + } + }); } diff --git a/backend/src/api/model/event_manager.rs b/backend/src/api/model/event_manager.rs index c1d919249..843e9a250 100644 --- a/backend/src/api/model/event_manager.rs +++ b/backend/src/api/model/event_manager.rs @@ -15,6 +15,8 @@ pub enum EventMessage { LibraryScanProgress(LibraryScanSummary), // Triggered when MetadataUpdateManager queue drains for an input InputMetadataUpdatesCompleted(Arc), + // Triggered when MetadataUpdateManager starts a new processing cycle for an input + InputMetadataUpdatesStarted(Arc), } pub struct EventManager { diff --git a/backend/src/api/model/metadata_update_manager.rs b/backend/src/api/model/metadata_update_manager.rs index 82831a6de..0135ae6e5 100644 --- a/backend/src/api/model/metadata_update_manager.rs +++ b/backend/src/api/model/metadata_update_manager.rs @@ -373,6 +373,11 @@ impl InputWorker { last_queue_log_at = Instant::now() .checked_sub(QUEUE_LOG_INTERVAL + Duration::from_secs(1)) .unwrap_or_else(Instant::now); + if let Some(app_state) = app_state_weak.as_ref().and_then(Weak::upgrade) { + app_state + .event_manager + .send_event(EventMessage::InputMetadataUpdatesStarted(input_name.clone())); + } debug!("Background metadata update queue has entries for input {input_name}; starting processing"); } diff --git a/backend/src/api/model/update_guard.rs b/backend/src/api/model/update_guard.rs index 48da1ef20..dbe28f067 100644 --- a/backend/src/api/model/update_guard.rs +++ b/backend/src/api/model/update_guard.rs @@ -29,6 +29,15 @@ impl UpdateGuard { .map(|permit| UpdateGuardPermit { _permit: permit }) } + pub async fn acquire_playlist_lock(&self) -> Option { + self.playlist + .clone() + .acquire_owned() + .await + .ok() + .map(|permit| UpdateGuardPermit { _permit: permit }) + } + pub fn try_library(&self) -> Option { self.library .clone() diff --git a/backend/src/api/scheduler.rs b/backend/src/api/scheduler.rs index cfe39d0ad..5d931adf0 100644 --- a/backend/src/api/scheduler.rs +++ b/backend/src/api/scheduler.rs @@ -67,7 +67,7 @@ async fn start_scheduler(client: reqwest::Client, expression: &str, app_state: A let metadata_manager = Arc::clone(&app_state.metadata_manager); sync_panel_api_exp_dates_on_boot(&app_state).await; exec_processing(&client, app_config, Arc::clone(&targets), Some(event_manager), - Some(playlist_state), Some(app_state.update_guard.clone()), disabled_headers, Some(provider_manager), Some(metadata_manager)).await; + Some(playlist_state), Some(app_state.update_guard.clone()), disabled_headers, Some(provider_manager), Some(metadata_manager), None, None).await; } () = cancel.cancelled() => { break; diff --git a/backend/src/main.rs b/backend/src/main.rs index bbfc17d19..c36a98c28 100644 --- a/backend/src/main.rs +++ b/backend/src/main.rs @@ -246,7 +246,7 @@ async fn start_in_cli_mode(cfg: Arc, targets: Arc) { reqwest::Client::new() }); // In CLI mode, we don't start background managers for events or providers - exec_processing(&client, cfg, targets, None, None, None, None, None, None).await; + exec_processing(&client, cfg, targets, None, None, None, None, None, 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 cca8af426..24806874c 100644 --- a/backend/src/processing/processor/playlist.rs +++ b/backend/src/processing/processor/playlist.rs @@ -14,6 +14,7 @@ use crate::api::model::{ ActiveProviderManager, EventManager, EventMessage, MetadataUpdateManager, PlaylistStorageState, ProviderIdType, ResolveReason, ResolveReasonSet, UpdateGuard, UpdateTask, }; + use crate::messaging::send_message; use crate::model::FetchedPlaylist; @@ -43,7 +44,7 @@ use shared::foundation::{get_field_value, set_field_value, ValueAccessor, ValueP use shared::model::xtream_const::XTREAM_CLUSTER; use shared::model::{ CounterModifier, FieldGetAccessor, FieldSetAccessor, InputType, ItemField, PlaylistGroup, PlaylistItem, - PlaylistItemType, PlaylistUpdateState, ProcessingOrder, StreamProperties, XtreamCluster, + PlaylistItemType, ProcessingOrder, StreamProperties, XtreamCluster, }; use shared::model::{InputStats, PlaylistStats, SourceStats, TargetStats, UUIDType}; use shared::utils::{ @@ -569,6 +570,11 @@ async fn download_input( if already_processed { // Use empty results, will load from disk below PlaylistDownloadResult::new(vec![], vec![], true, false) + } else if ctx.pre_processed_inputs.as_ref().is_some_and(|s| s.contains(&input.name)) { + // Input was already processed in a prior session; skip download and mark as processed + // for this session to avoid redundant lock acquisitions from other targets. + ctx.mark_input_downloaded(input.name.clone()).await; + PlaylistDownloadResult::new(vec![], vec![], true, false) } else { let res = playlist_download_from_input(&ctx.client, &ctx.config, input).await; // Mark as processed if NO critical errors? @@ -640,6 +646,7 @@ pub struct PlaylistProcessingContext { // New field for STRM probes & background updates pub provider_manager: Option>, pub metadata_manager: Option>, + pub pre_processed_inputs: Option>>>, } impl PlaylistProcessingContext { @@ -1192,14 +1199,19 @@ pub async fn exec_processing( disabled_headers: Option, provider_manager: Option>, metadata_manager: Option>, + pre_processed_inputs: Option>>, + acquired_permit: Option, ) { - let _guard = if let Some(guard) = update_guard { - if let Some(permit) = guard.try_playlist() { + let _guard = if let Some(permit) = acquired_permit { + Some(permit) + } else if let Some(guard) = &update_guard { + let lock_result: Option = guard.try_playlist(); + if let Some(permit) = lock_result { 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)); + if let Some(events) = event_manager.as_deref() { + events.send_event(EventMessage::PlaylistUpdate(shared::model::PlaylistUpdateState::Failure)); } return; } @@ -1221,6 +1233,7 @@ pub async fn exec_processing( disabled_headers, provider_manager, metadata_manager, + pre_processed_inputs: pre_processed_inputs.map(Arc::new), }; let start_time = Instant::now(); @@ -1242,18 +1255,18 @@ pub async fn exec_processing( // send errors if let Some(message) = get_errors_notify_message!(errors, 255) { - if let Some(events) = &event_manager { - events.send_event(EventMessage::PlaylistUpdate(PlaylistUpdateState::Failure)); + if let Some(events) = event_manager.as_deref() { + events.send_event(EventMessage::PlaylistUpdate(shared::model::PlaylistUpdateState::Failure)); } send_message(&app_config, client, MessageContent::event_error(message)).await; - } else if let Some(events) = &event_manager { - events.send_event(EventMessage::PlaylistUpdate(PlaylistUpdateState::Success)); + } else if let Some(events) = event_manager.as_deref() { + events.send_event(EventMessage::PlaylistUpdate(shared::model::PlaylistUpdateState::Success)); } let elapsed = start_time.elapsed().as_secs(); let update_finished_message = format!("🌷 Update process finished! Took {elapsed} secs."); - if let Some(events) = &event_manager { + if let Some(events) = event_manager.as_deref() { events.send_event(EventMessage::PlaylistUpdateProgress( "Playlist Update".to_string(), update_finished_message.clone(),