diff --git a/backend/src/api/endpoints/api_playlist_utils.rs b/backend/src/api/endpoints/api_playlist_utils.rs index d594411d1..cb1d1031e 100644 --- a/backend/src/api/endpoints/api_playlist_utils.rs +++ b/backend/src/api/endpoints/api_playlist_utils.rs @@ -1,4 +1,4 @@ -use crate::model::{AppConfig, Config, ConfigInput, ConfigTarget}; +use crate::model::{AppConfig, ConfigInput, ConfigTarget}; use crate::repository::{m3u_repository, xtream_repository}; use crate::utils::{m3u, xtream}; use crate::utils; @@ -144,14 +144,15 @@ pub(in crate::api::endpoints) async fn get_playlist_for_target(cfg_target: Optio (axum::http::StatusCode::BAD_REQUEST, axum::Json(json!({"error": "Invalid Arguments"}))).into_response() } -pub(in crate::api::endpoints) async fn get_playlist(client: &reqwest::Client, cfg_input: Option<&Arc>, cfg: &Arc, accept: Option<&str>) -> impl IntoResponse + Send { +pub(in crate::api::endpoints) async fn get_playlist(client: &reqwest::Client, cfg_input: Option<&Arc>, app_config: &Arc, accept: Option<&str>) -> impl IntoResponse + Send { + let cfg = app_config.config.load(); match cfg_input { Some(input) => { let (result, errors) = match input.input_type { - InputType::M3u | InputType::M3uBatch => m3u::download_m3u_playlist(client, cfg, input).await, + InputType::M3u | InputType::M3uBatch => m3u::download_m3u_playlist(client, &cfg, input).await, InputType::Xtream | InputType::XtreamBatch => { - let (pl, err, _) = xtream::download_xtream_playlist(cfg, client, input, None).await; + let (pl, err, _) = xtream::download_xtream_playlist(app_config, client, input, None).await; (pl, err) } InputType::Library => { diff --git a/backend/src/api/endpoints/v1_api_playlist.rs b/backend/src/api/endpoints/v1_api_playlist.rs index 50503dd76..1b3b4d528 100644 --- a/backend/src/api/endpoints/v1_api_playlist.rs +++ b/backend/src/api/endpoints/v1_api_playlist.rs @@ -85,20 +85,20 @@ async fn playlist_content( axum::extract::State(app_state): axum::extract::State>, axum::extract::Json(playlist_req): axum::extract::Json, ) -> impl IntoResponse + Send { - let config = app_state.app_config.config.load(); + let _config = app_state.app_config.config.load(); let client = app_state.http_client.load(); match playlist_req { PlaylistRequest::Target(target_id) => { get_playlist_for_target(app_state.app_config.get_target_by_id(target_id).as_deref(), &app_state.app_config, accept.as_deref()).await.into_response() } PlaylistRequest::Input(input_id) => { - get_playlist(client.as_ref(), app_state.app_config.get_input_by_id(input_id).as_ref(), &config, accept.as_deref()).await.into_response() + get_playlist(client.as_ref(), app_state.app_config.get_input_by_id(input_id).as_ref(), &app_state.app_config, accept.as_deref()).await.into_response() } PlaylistRequest::CustomXtream(xtream) => { match Url::parse(&xtream.url) { Ok(parsed) if parsed.scheme() == "http" || parsed.scheme() == "https" => { let input = Arc::new(create_config_input_for_xtream(&xtream.username, &xtream.password, &xtream.url)); - get_playlist(client.as_ref(), Some(&input), &config, accept.as_deref()).await.into_response() + get_playlist(client.as_ref(), Some(&input), &app_state.app_config, accept.as_deref()).await.into_response() } _ => { (axum::http::StatusCode::BAD_REQUEST, axum::Json(json!({"error": "Invalid url scheme; only http/https are allowed"}))).into_response() @@ -109,7 +109,7 @@ async fn playlist_content( match Url::parse(&m3u.url) { Ok(parsed) if parsed.scheme() == "http" || parsed.scheme() == "https" => { let input = Arc::new(create_config_input_for_m3u(&m3u.url)); - get_playlist(client.as_ref(), Some(&input), &config, accept.as_deref()).await.into_response() + get_playlist(client.as_ref(), Some(&input), &app_state.app_config, accept.as_deref()).await.into_response() } _ => { (axum::http::StatusCode::BAD_REQUEST, axum::Json(json!({"error": "Invalid url scheme; only http/https are allowed"}))).into_response() diff --git a/backend/src/processing/processor/playlist.rs b/backend/src/processing/processor/playlist.rs index 7eba638f9..f216be126 100644 --- a/backend/src/processing/processor/playlist.rs +++ b/backend/src/processing/processor/playlist.rs @@ -319,7 +319,7 @@ async fn playlist_download_from_input(client: &reqwest::Client, app_config: &Arc let (p, e) = m3u::download_m3u_playlist(client, config, input).await; (p, e, false) } - InputType::Xtream => xtream::download_xtream_playlist(config, client, input, clusters_to_download.as_deref()).await, + InputType::Xtream => xtream::download_xtream_playlist(app_config, client, input, clusters_to_download.as_deref()).await, InputType::M3uBatch | InputType::XtreamBatch => (vec![], vec![], false), InputType::Library => { let (p, e) = library::download_library_playlist(client, app_config, input).await; @@ -398,14 +398,13 @@ async fn process_source(source_idx: usize, ctx: &PlaylistProcessingContext) -> ( (playlist, download_err) }; - let (tvguide, mut tvguide_errors) = if input.input_type == InputType::Library { - (None, vec!()) + let tvguide = if input.input_type == InputType::Library { + None } else { - download_input_epg(ctx,input, &mut errors).await + download_input_epg(ctx,input, &mut error_list).await }; errors.append(&mut error_list); - errors.append(&mut tvguide_errors); let group_count = playlist_groups.get_group_count(); let channel_count = playlist_groups.get_channel_count(); let input_name = &input.name; @@ -455,16 +454,16 @@ async fn process_source(source_idx: usize, ctx: &PlaylistProcessingContext) -> ( } async fn download_input_epg(ctx: &PlaylistProcessingContext, input: &Arc, - error_list: &mut [TuliproxError]) -> (Option, Vec) { + error_list: &mut Vec) -> Option { // Download epg for input - let (tvguide, tvguide_errors) = if error_list.is_empty() { - debug!("Downloading epg for input '{}'", input.name); + let (tvguide, mut tvguide_errors) = if error_list.is_empty() { let working_dir = &ctx.config.config.load().working_dir; epg::get_xmltv(ctx, input, working_dir).await } else { (None, vec![]) }; - (tvguide, tvguide_errors) + error_list.append(&mut tvguide_errors); + tvguide } async fn download_input(ctx: &PlaylistProcessingContext, input: &Arc) diff --git a/backend/src/utils/network/epg.rs b/backend/src/utils/network/epg.rs index 6bed93d54..62d1106c6 100644 --- a/backend/src/utils/network/epg.rs +++ b/backend/src/utils/network/epg.rs @@ -46,12 +46,12 @@ async fn download_epg_file(url: &str, ctx: &PlaylistProcessingContext, input: &C } let lock_key = persist_file_path.display().to_string(); - let _input_lock = ctx.get_input_lock(&lock_key); + let _input_lock = ctx.get_input_lock(&lock_key).await; if ctx.is_input_downloaded(&lock_key).await { return Ok(persist_file_path); } - + debug!("Downloading epg for input '{}'", input.name); match request::get_input_epg_content_as_file(&ctx.client, input, working_dir, url, &persist_file_path).await { Ok(path) => { ctx.mark_input_downloaded(lock_key.clone()).await; diff --git a/backend/src/utils/network/xtream.rs b/backend/src/utils/network/xtream.rs index fe0dc56d0..6c2f18886 100644 --- a/backend/src/utils/network/xtream.rs +++ b/backend/src/utils/network/xtream.rs @@ -1,6 +1,6 @@ use crate::api::model::AppState; use crate::messaging::send_message; -use crate::model::{is_input_expired, xtream_mapping_option_from_target_options, Config, ConfigInput, ConfigTarget, XtreamLoginInfo, XtreamTargetOutput}; +use crate::model::{is_input_expired, xtream_mapping_option_from_target_options, AppConfig, Config, ConfigInput, ConfigTarget, XtreamLoginInfo, XtreamTargetOutput}; use crate::model::{InputSource, ProxyUserCredentials}; use crate::processing::parser::xtream; use crate::processing::parser::xtream::parse_xtream_series_info; @@ -336,8 +336,9 @@ pub async fn notify_account_expire(exp_date: Option, cfg: &Config, client: } } -pub async fn download_xtream_playlist(cfg: &Arc, client: &reqwest::Client, input: &ConfigInput, clusters: Option<&[XtreamCluster]>) +pub async fn download_xtream_playlist(app_config: &Arc, client: &reqwest::Client, input: &ConfigInput, clusters: Option<&[XtreamCluster]>) -> (Vec, Vec, bool) { + let cfg = app_config.config.load(); let input_source: InputSource = { match input.staged.as_ref() { None => input.into(), @@ -351,9 +352,9 @@ pub async fn download_xtream_playlist(cfg: &Arc, client: &reqwest::Clien let base_url = get_xtream_stream_url_base(&input_source.url, username, password); let input_source_login = input_source.with_url(base_url.clone()); - check_alias_user_state(cfg, client, input).await; + check_alias_user_state(&cfg, client, input).await; - if let Err(err) = xtream_login(cfg, client, &input_source_login, username).await { + if let Err(err) = xtream_login(&cfg, client, &input_source_login, username).await { error!("Could not log in with xtream user {username} for provider {}. {err}", input.name); return (Vec::with_capacity(0), vec![err], false); } @@ -381,7 +382,7 @@ pub async fn download_xtream_playlist(cfg: &Arc, client: &reqwest::Clien (Ok(category_content), Ok(stream_content)) => { if cfg.disk_based_processing { // trace!("Using disk input playlist optimization for cluster {}", xtream_cluster); - if let Err(err) = process_xtream_cluster_to_disk(cfg, input, *xtream_cluster, category_content, stream_content).await { + if let Err(err) = process_xtream_cluster_to_disk(app_config, input, *xtream_cluster, category_content, stream_content).await { error!("process_xtream_cluster_to_disk failed: {err}"); errors.push(err); } else { @@ -498,15 +499,16 @@ pub fn create_vod_info_from_item(target: &ConfigTarget, user: &ProxyUserCredenti const BATCH_SIZE: usize = 1000; async fn process_xtream_cluster_to_disk( - cfg: &Arc, + app_config: &Arc, input: &ConfigInput, cluster: XtreamCluster, categories: DynReader, streams: DynReader, ) -> Result<(), TuliproxError> { + let cfg = app_config.config.load(); // trace!("Starting process_xtream_cluster_to_disk for cluster {}", cluster); let storage_path = { - ensure_input_storage_path(cfg, &input.name)? + ensure_input_storage_path(&cfg, &input.name)? }; let xtream_path = xtream_get_file_path(&storage_path, cluster); @@ -587,6 +589,9 @@ async fn process_xtream_cluster_to_disk( // 2. Success! Swap temporary files to permanent ones let tmp_xtream_path = xtream_path.with_extension("tmp"); + // Acquire write lock to serialize compact/swap/cleanup operations across concurrent API calls + let swap_lock = app_config.file_locks.write_lock(&xtream_path).await; + if let Ok(mut tree_update) = BPlusTreeUpdate::::try_new(&tmp_xtream_path) { // Compact the TEMPORARY file (tmp_xtream_path) in place. // We ensure the .tmp file is compacted before we rename it to the final destination, @@ -617,6 +622,7 @@ async fn process_xtream_cluster_to_disk( let _ = tokio::fs::remove_file(tmp_col_path).await; } + drop(swap_lock); // trace!("Cluster {} updated successfully", cluster); Ok(()) } diff --git a/frontend/src/app/components/config/library_config_view.rs b/frontend/src/app/components/config/library_config_view.rs index 8ada299d7..1ad2078b7 100644 --- a/frontend/src/app/components/config/library_config_view.rs +++ b/frontend/src/app/components/config/library_config_view.rs @@ -389,9 +389,7 @@ pub fn LibraryConfigView() -> Html { <>
{ edit_field_bool!(form_state, translate.t(LABEL_ENABLED), enabled, LibraryConfigFormAction::Enabled) } - { render_scan_directories_edit() } -