mod playlist; mod stalker; mod stalker_refresh; mod xtream; // mod affix; mod deduplicate; mod epg; mod library; mod probe_handle_guard; mod providers; mod resolve_options; mod sort; mod stream_probe; mod xtream_series; mod xtream_vod; pub use self::{ epg::*, playlist::*, probe_handle_guard::*, providers::*, resolve_options::*, stalker::{download_stalker_playlist, re_resolve_stalker_url, StalkerCluster}, stalker_refresh::StalkerRefreshMode, stream_probe::*, xtream::*, xtream_series::*, xtream_vod::*, }; use tokio_util::sync::CancellationToken; use tuliprox_core::model::ProviderHandle; pub const FOREGROUND_BATCH_SIZE: usize = 200; pub const FOREGROUND_RETRY_BATCH_MAX_SIZE: usize = FOREGROUND_BATCH_SIZE * 4; pub const FOREGROUND_MIN_RETRY_DELAY_SECS: u64 = 1; pub fn select_cancel_token<'a>( acquired_handle: Option<&'a ProviderHandle>, active_handle: Option<&'a ProviderHandle>, ) -> Option<&'a CancellationToken> { acquired_handle .and_then(|h| h.cancel_token.as_ref()) .or_else(|| active_handle.and_then(|h| h.cancel_token.as_ref())) } // // fn get_resolve__options(target: &ConfigTarget, fpl: &FetchedPlaylist) -> bool // #[macro_export] macro_rules! create_resolve_options_function_for_xtream_target { ($cluster:ident) => { paste::paste! { fn [](target: &ConfigTarget, fpl: &FetchedPlaylist) -> $crate::processor::ResolveOptions { match target.get_xtream_output() { Some(_) => { let input_options = fpl.input.options.as_ref(); let input_is_xtream = fpl.input.input_type.is_xtream(); let ( resolve_tmdb_missing, input_resolve_enabled, input_probe_enabled, input_resolve_delay, resolve_background ) = if let Some(options) = input_options { ( options.has_flag(tuliprox_core::model::ConfigInputFlags::ResolveTmdb), options.has_flag(tuliprox_core::model::ConfigInputFlags::[]), options.has_flag(tuliprox_core::model::ConfigInputFlags::[]), options.resolve_delay, options.has_flag(tuliprox_core::model::ConfigInputFlags::ResolveBackground) ) } else { ( false, false, false, shared::defaults::default_resolve_delay_secs(), shared::defaults::default_resolve_background(), ) }; let resolve_enabled = input_resolve_enabled; let probe_enabled = input_probe_enabled; let resolve_delay = input_resolve_delay; let mut flags = $crate::processor::ResolveOptionsFlagsSet::new(); if resolve_enabled && input_is_xtream { flags.set($crate::processor::ResolveOptionsFlags::Resolve); } if resolve_tmdb_missing { flags.set($crate::processor::ResolveOptionsFlags::TmdbMissing); } if input_is_xtream && probe_enabled { flags.set($crate::processor::ResolveOptionsFlags::Probe); } if resolve_background { flags.set($crate::processor::ResolveOptionsFlags::Background); } $crate::processor::ResolveOptions { flags, resolve_delay, } }, None => $crate::processor::ResolveOptions::default(), } } } }; } use create_resolve_options_function_for_xtream_target; /// Foreground retry helper that retries each queued item at most once. /// /// `retry_delay_secs` is applied sequentially per item, so the total wall-clock /// delay is roughly `retry_delay_secs * retry_item_count` plus network/DB time. macro_rules! process_foreground_retry_once { ( ctx: $ctx:expr, fpl: $fpl:expr, filter: $filter:expr, retry_once_ids: $retry_once_ids:ident, retry_delay_secs: $retry_delay_secs:expr, xtream_path: $xtream_path:ident, db_query_holder: $db_query_holder:ident, db_lock_holder: $db_lock_holder:ident, batch: $batch:ident, batch_size: $batch_size:expr, retry_batch_max_len: $retry_batch_max_len:expr, processed_count: $processed_count:ident, query_error_context: $query_error_context:expr, reasons: |$pli_reasons:ident| $reasons_expr:expr, update: |$active_provider:ident, $pli_update:ident, $provider_id:ident, $reasons:ident, $db_query_ref:ident| $update_expr:expr, apply_properties: |$pli_apply:ident, $updated_props:ident| $apply_expr:expr, persist: |$updates:ident| $persist_expr:expr, on_persist_error: |$persist_err:ident| $on_persist_error_expr:expr, on_retry_error: |$pli_error:ident, $retry_err:ident| $on_retry_error_expr:expr, on_after_attempt: |$pli_after:ident, $retry_succeeded:ident| $on_after_attempt_expr:expr $(,)? ) => { for __pli in $fpl.items_mut() { if !($filter)(__pli) { continue; } let __provider_id = if let Ok(__uid) = __pli.header.id.parse::() { tuliprox_core::model::ProviderIdType::Id(__uid) } else { tuliprox_core::model::ProviderIdType::from(&*__pli.header.id) }; if !$retry_once_ids.remove(&__provider_id) { continue; } let mut __retry_succeeded = false; let __reasons = { let $pli_reasons = &mut *__pli; $reasons_expr }; if !__reasons.is_empty() { if let Some(__active_provider) = $ctx.provider_manager.as_ref() { // Do not hold a read lock over the retry delay window. if $db_query_holder.is_some() { $db_query_holder = None; $db_lock_holder = None; } tokio::time::sleep(std::time::Duration::from_secs($retry_delay_secs)).await; if $db_query_holder.is_none() && $xtream_path.exists() { let __file_lock = $ctx.config.file_locks.read_lock(&$xtream_path).await; let __xtream_path = $xtream_path.clone(); let __query = match tokio::task::spawn_blocking(move || { tuliprox_repository::BPlusTreeQuery::::try_new( &__xtream_path, ) }) .await { Ok(Ok(__query)) => Some((__query, __file_lock)), Ok(Err(__err)) => { log::error!("Failed to open BPlusTreeQuery for {}: {__err}", $query_error_context); None } Err(__err) => { log::error!("Failed to open BPlusTreeQuery for {}: {__err}", $query_error_context); None } }; if let Some((__query, __guard)) = __query { $db_query_holder = Some(std::sync::Arc::new(parking_lot::Mutex::new(__query))); $db_lock_holder = Some(__guard); } } let __db_query_ref = $db_query_holder.as_ref().map(std::sync::Arc::clone); let __update_future = { let $active_provider = __active_provider; let $pli_update = &mut *__pli; let $provider_id = __provider_id.clone(); let $reasons = &__reasons; let $db_query_ref = __db_query_ref; $update_expr }; match __update_future.await { Ok(Some(__updated_props)) => { { let $pli_apply = &mut *__pli; let $updated_props = &__updated_props; $apply_expr } $batch.push((__provider_id.clone(), __updated_props)); if $batch.len() >= $batch_size { $db_query_holder = None; $db_lock_holder = None; let __updates: Vec<(u32, _)> = $batch .iter() .filter_map(|(__id, __props)| { // Foreground retry batches can include text provider IDs (e.g. M3U). // Persist batch functions for these paths are keyed by numeric Xtream IDs, // so text IDs are intentionally skipped here. if let tuliprox_core::model::ProviderIdType::Id(__vid) = __id { Some((*__vid, __props.clone())) } else { None } }) .collect(); if __updates.is_empty() { $batch.clear(); } else { let __persist_future = { let $updates = __updates; $persist_expr }; match __persist_future.await { Ok(()) => $batch.clear(), Err($persist_err) => { $on_persist_error_expr; if $batch.len() > $retry_batch_max_len { let __drop_count = $batch.len().saturating_sub($retry_batch_max_len); if __drop_count > 0 { $batch.drain(0..__drop_count); } } } } } } $processed_count += 1; __retry_succeeded = true; } Ok(None) => {} Err($retry_err) => { let $pli_error = &*__pli; $on_retry_error_expr; } } } } { let $pli_after = &mut *__pli; let $retry_succeeded = __retry_succeeded; $on_after_attempt_expr; } } }; } // `macro_rules!` without `#[macro_export]`: the re-export cannot be wider than // the crate. pub(crate) use process_foreground_retry_once;