diff --git a/README.md b/README.md index d544ac974..103d1dacd 100644 --- a/README.md +++ b/README.md @@ -1465,8 +1465,24 @@ Now you can do `nginx` configuration like proxy_ssl_session_reuse off; proxy_set_header Host $http_host; proxy_redirect off; + proxy_buffering off; + proxy_request_buffering off; + proxy_cache off; + tcp_nopush on; + tcp_nodelay on; } ``` +When you use nginx be sure to have +``` + proxy_redirect off; + proxy_buffering off; + proxy_request_buffering off; + proxy_cache off; + tcp_nopush on; + tcp_nodelay on; +``` +because without this config you could get very high cpu peaks. + You can also use traefik as reverse proxy server in front of your tuliprox instance. However if you wan't to use paths, you must note that the path for web-ui and api-proxy must be different. In this short example used paths are: * web-ui: tuliprox * api-proxy: tv diff --git a/backend/src/api/endpoints/custom_video_stream_api.rs b/backend/src/api/endpoints/custom_video_stream_api.rs index 125891bf8..8761a1e97 100644 --- a/backend/src/api/endpoints/custom_video_stream_api.rs +++ b/backend/src/api/endpoints/custom_video_stream_api.rs @@ -12,7 +12,9 @@ async fn cvs_api( axum::extract::State(app_state): axum::extract::State>, ) -> impl IntoResponse + Send { - let Ok(custom_video_type) = CustomVideoStreamType::from_str(&stream_type) else { + let cvs_type = stream_type.strip_suffix(".ts").unwrap_or(&stream_type); + + let Ok(custom_video_type) = CustomVideoStreamType::from_str(cvs_type) else { return axum::http::StatusCode::NOT_FOUND.into_response(); }; diff --git a/backend/src/api/endpoints/hls_api.rs b/backend/src/api/endpoints/hls_api.rs index 6208b8cdb..fb18297a5 100644 --- a/backend/src/api/endpoints/hls_api.rs +++ b/backend/src/api/endpoints/hls_api.rs @@ -129,7 +129,7 @@ pub(in crate::api) async fn handle_hls_stream_request( let custom_stream_response = app_state.app_config.custom_stream_response.load(); if custom_stream_response.as_ref().and_then(|c| c.channel_unavailable.as_ref()).is_some() { let url = format!( - "{}/{CUSTOM_VIDEO_PREFIX}/{}/{}/{}", + "{}/{CUSTOM_VIDEO_PREFIX}/{}/{}/{}.ts", &server_info.get_base_url(), user.username, user.password, diff --git a/backend/src/api/endpoints/m3u_api.rs b/backend/src/api/endpoints/m3u_api.rs index 68c2df411..d950206a2 100644 --- a/backend/src/api/endpoints/m3u_api.rs +++ b/backend/src/api/endpoints/m3u_api.rs @@ -113,7 +113,7 @@ async fn m3u_api_stream( let (action_stream_id, stream_ext) = separate_number_and_remainder(stream_req.stream_id); let virtual_id: u32 = try_result_bad_request!(action_stream_id.trim().parse()); let pli = try_result_not_found!( - m3u_get_item_for_stream_id(virtual_id, &app_state, &target).await, + m3u_get_item_for_stream_id(virtual_id, app_state, &target).await, true, format!("Failed to read m3u item for stream id {}", virtual_id) ); diff --git a/backend/src/api/endpoints/v1_api_playlist.rs b/backend/src/api/endpoints/v1_api_playlist.rs index 0ad1fb8fb..e040f5263 100644 --- a/backend/src/api/endpoints/v1_api_playlist.rs +++ b/backend/src/api/endpoints/v1_api_playlist.rs @@ -62,7 +62,8 @@ async fn playlist_update( Ok(valid_targets) => { let app_config = Arc::clone(&app_state.app_config); let event_manager = Arc::clone(&app_state.event_manager); - tokio::spawn(playlist::exec_processing(Arc::clone(&app_state.http_client.load()), app_config, Arc::new(valid_targets), Some(event_manager))); + let playlist_state = Arc::clone(&app_state.playlists); + tokio::spawn(playlist::exec_processing(Arc::clone(&app_state.http_client.load()), app_config, Arc::new(valid_targets), Some(event_manager), Some(playlist_state))); axum::http::StatusCode::ACCEPTED.into_response() } Err(err) => { diff --git a/backend/src/api/endpoints/xtream_api.rs b/backend/src/api/endpoints/xtream_api.rs index 6a0d3118d..814a6b4b1 100644 --- a/backend/src/api/endpoints/xtream_api.rs +++ b/backend/src/api/endpoints/xtream_api.rs @@ -249,7 +249,7 @@ async fn xtream_player_api_stream( let (pli, mapping) = try_result_not_found!( xtream_repository::xtream_get_item_for_stream_id( virtual_id, - &app_state, + app_state, &target, None ).await, @@ -426,7 +426,7 @@ async fn xtream_player_api_stream_with_token( let (pli, _mapping) = try_result_bad_request!( xtream_repository::xtream_get_item_for_stream_id( virtual_id, - &app_state, + app_state, &target, None ).await, @@ -729,7 +729,7 @@ async fn xtream_player_api_resource( let (pli, _) = try_result_bad_request!( xtream_repository::xtream_get_item_for_stream_id( virtual_id, - &app_state, + app_state, &target, None ).await, @@ -969,7 +969,7 @@ async fn xtream_get_stream_info_response( if let Ok((pli, virtual_record)) = xtream_repository::xtream_get_item_for_stream_id( virtual_id, - &app_state, + app_state, target, Some(cluster), ).await { @@ -1054,7 +1054,7 @@ async fn xtream_get_short_epg( if let Ok((pli, _)) = xtream_repository::xtream_get_item_for_stream_id( virtual_id, - &app_state, + app_state, target, None, ).await { @@ -1215,7 +1215,7 @@ async fn xtream_get_catchup_response( let virtual_id: u32 = try_result_bad_request!(FromStr::from_str(stream_id)); let (pli, _) = try_result_bad_request!(xtream_repository::xtream_get_item_for_stream_id( virtual_id, - &app_state, + app_state, target, Some(XtreamCluster::Live) ).await); diff --git a/backend/src/api/main_api.rs b/backend/src/api/main_api.rs index 0967fe69c..6888fcade 100644 --- a/backend/src/api/main_api.rs +++ b/backend/src/api/main_api.rs @@ -1,4 +1,3 @@ -use std::collections::HashMap; use crate::api::api_utils::{get_build_time, get_server_time}; use crate::api::config_watch::exec_config_watch; use crate::api::endpoints::hdhomerun_api::hdhr_api_register; @@ -9,7 +8,7 @@ use crate::api::endpoints::web_index::{index_register_with_path, index_register_ use crate::api::endpoints::websocket_api::ws_api_register; use crate::api::endpoints::xmltv_api::xmltv_api_register; use crate::api::endpoints::xtream_api::xtream_api_register; -use crate::api::model::ActiveProviderManager; +use crate::api::model::{ActiveProviderManager, PlaylistStorageState}; use crate::api::model::ActiveUserManager; use crate::api::model::DownloadQueue; use crate::api::model::EventManager; @@ -29,7 +28,6 @@ use log::{error, info}; use std::io::ErrorKind; use std::path::PathBuf; use std::sync::{Arc}; -use tokio::sync::RwLock; use tokio_util::sync::CancellationToken; use tower_governor::key_extractor::SmartIpKeyExtractor; use shared::utils::{concat_path_leading_slash}; @@ -94,7 +92,7 @@ fn create_shared_data( active_provider, event_manager, cancel_tokens: Arc::new(ArcSwap::from_pointee(CancelTokens::default())), - playlists: Arc::new(RwLock::new(HashMap::new())), + playlists: Arc::new(PlaylistStorageState::new()), } } @@ -111,8 +109,9 @@ fn exec_update_on_boot( if update_on_boot { let app_state_clone = Arc::clone(&app_state.app_config); let targets_clone = Arc::clone(targets); + let playlist_state = Arc::clone(&app_state.playlists); tokio::spawn( - async move { playlist::exec_processing(client, app_state_clone, targets_clone, None).await }, + async move { playlist::exec_processing(client, app_state_clone, targets_clone, None, Some(playlist_state)).await }, ); } } diff --git a/backend/src/api/model/app_state.rs b/backend/src/api/model/app_state.rs index aad3ea61d..9cade932c 100644 --- a/backend/src/api/model/app_state.rs +++ b/backend/src/api/model/app_state.rs @@ -172,8 +172,8 @@ pub struct PlaylistXtreamStorage { pub type PlaylistM3uStorage = BPlusTree; pub enum PlaylistStorage { - M3uPlaylist(PlaylistM3uStorage), - XtreamPlaylist(PlaylistXtreamStorage), + M3uPlaylist(Box), + XtreamPlaylist(Box), } pub struct TargetPlaylistStorage { @@ -181,7 +181,53 @@ pub struct TargetPlaylistStorage { pub m3u: Option, } -type TargetPlaylistStorageMap = HashMap; +pub type TargetPlaylistStorageMap = HashMap; + +pub struct PlaylistStorageState { + pub data: RwLock, +} + +impl PlaylistStorageState { + + pub(crate) fn new() -> Self { + Self { + data: RwLock::new(HashMap::new()), + } + } + + pub async fn cache_playlist(&self, target_name: &str, playlist: PlaylistStorage) { + match playlist { + PlaylistStorage::M3uPlaylist(m3u_playlist) => { + match self.data.write().await.entry(target_name.to_string()) { + std::collections::hash_map::Entry::Occupied(mut entry) => { + let storage = entry.get_mut(); + storage.m3u = Some(*m3u_playlist); + } + std::collections::hash_map::Entry::Vacant(entry) => { + entry.insert(TargetPlaylistStorage { + xtream: None, + m3u: Some(*m3u_playlist), + }); + } + } + } + PlaylistStorage::XtreamPlaylist(xtream_playlist) => { + match self.data.write().await.entry(target_name.to_string()) { + std::collections::hash_map::Entry::Occupied(mut entry) => { + let storage = entry.get_mut(); + storage.xtream = Some(*xtream_playlist); + } + std::collections::hash_map::Entry::Vacant(entry) => { + entry.insert(TargetPlaylistStorage { + xtream: Some(*xtream_playlist), + m3u: None, + }); + } + } + } + } + } +} #[derive(Clone)] pub struct AppState { @@ -195,7 +241,7 @@ pub struct AppState { pub active_provider: Arc, pub event_manager: Arc, pub cancel_tokens: Arc>, - pub playlists: Arc> + pub playlists: Arc, } impl AppState { @@ -277,36 +323,7 @@ impl AppState { } pub async fn cache_playlist(&self, target_name: &str, playlist: PlaylistStorage) { - match playlist { - PlaylistStorage::M3uPlaylist(m3u_playlist) => { - match self.playlists.write().await.entry(target_name.to_string()) { - std::collections::hash_map::Entry::Occupied(mut entry) => { - let storage = entry.get_mut(); - storage.m3u = Some(m3u_playlist); - } - std::collections::hash_map::Entry::Vacant(entry) => { - entry.insert(TargetPlaylistStorage { - xtream: None, - m3u: Some(m3u_playlist), - }); - } - } - } - PlaylistStorage::XtreamPlaylist(xtream_playlist) => { - match self.playlists.write().await.entry(target_name.to_string()) { - std::collections::hash_map::Entry::Occupied(mut entry) => { - let storage = entry.get_mut(); - storage.xtream = Some(xtream_playlist); - } - std::collections::hash_map::Entry::Vacant(entry) => { - entry.insert(TargetPlaylistStorage { - xtream: Some(xtream_playlist), - m3u: None, - }); - } - } - } - } + self.playlists.cache_playlist(target_name, playlist).await; } } diff --git a/backend/src/api/scheduler.rs b/backend/src/api/scheduler.rs index ba5afa8a1..d258696ca 100644 --- a/backend/src/api/scheduler.rs +++ b/backend/src/api/scheduler.rs @@ -59,7 +59,8 @@ async fn start_scheduler(client: Arc, expression: &str, app_sta () = tokio::time::sleep_until(tokio::time::Instant::from(datetime_to_instant(datetime))) => { let app_config = Arc::clone(&app_state.app_config); let event_manager = Arc::clone(&app_state.event_manager); - exec_processing(Arc::clone(&client), app_config, Arc::clone(&targets), Some(event_manager)).await; + let playlist_state = app_state.playlists.clone(); + exec_processing(Arc::clone(&client), app_config, Arc::clone(&targets), Some(event_manager), Some(playlist_state)).await; } () = cancel.cancelled() => { break; diff --git a/backend/src/main.rs b/backend/src/main.rs index bbcdaae25..efdf7c59c 100644 --- a/backend/src/main.rs +++ b/backend/src/main.rs @@ -167,7 +167,7 @@ async fn start_in_cli_mode(cfg: Arc, targets: Arc) { error!("Failed to build client {err}"); reqwest::Client::new() }); - playlist::exec_processing(Arc::new(client), cfg, targets, None).await; + playlist::exec_processing(Arc::new(client), cfg, targets, 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 8ca8c8337..b3e4c7d35 100644 --- a/backend/src/processing/processor/playlist.rs +++ b/backend/src/processing/processor/playlist.rs @@ -34,7 +34,7 @@ use shared::model::{CounterModifier, FieldGetAccessor, FieldSetAccessor, InputTy PlaylistGroup, PlaylistItem, PlaylistUpdateState, ProcessingOrder, UUIDType, XtreamCluster}; use shared::utils::default_as_default; use std::time::Instant; -use crate::api::model::{EventManager, EventMessage}; +use crate::api::model::{EventManager, EventMessage, PlaylistStorageState}; fn is_valid(pli: &PlaylistItem, target: &ConfigTarget) -> bool { let provider = ValueProvider { pli }; @@ -244,8 +244,9 @@ async fn playlist_download_from_input(client: &Arc, config: &Ar } async fn process_source(client: Arc, cfg: Arc, source_idx: usize, - user_targets: Arc, event_manager: Option>) - -> (Vec, Vec, Vec) { + user_targets: Arc, event_manager: Option>, + playlist_state: Option<&Arc> +)-> (Vec, Vec, Vec) { let sources = cfg.sources.load(); let mut errors = vec![]; let mut input_stats = HashMap::::new(); @@ -301,7 +302,7 @@ async fn process_source(client: Arc, cfg: Arc, sourc for target in &source.targets { let event_manager_clone = event_manager_clone.clone(); if is_target_enabled(target, &user_targets) { - match process_playlist_for_target(&cfg, Arc::clone(&client), &mut source_playlists, target, &mut input_stats, &mut errors, event_manager_clone).await { + match process_playlist_for_target(&cfg, Arc::clone(&client), &mut source_playlists, target, &mut input_stats, &mut errors, event_manager_clone, playlist_state).await { Ok(()) => { target_stats.push(TargetStats::success(&target.name)); } @@ -335,7 +336,9 @@ fn create_input_stat(group_count: usize, channel_count: usize, error_count: usiz } } -async fn process_sources(client: Arc, config: &Arc, user_targets: Arc, event_manager: Option>) -> (Vec, Vec) { +async fn process_sources(client: Arc, config: &Arc, user_targets: Arc, + event_manager: Option>, playlist_state: Option<&Arc> +) -> (Vec, Vec) { let mut handle_list = vec![]; let thread_num = config.config.load().threads; let sources = config.sources.load(); @@ -360,6 +363,7 @@ async fn process_sources(client: Arc, config: &Arc, let event_manager = event_manager.clone(); if process_parallel { let http_client = Arc::clone(&client); + let playlist_state = playlist_state.cloned(); let handles = &mut handle_list; let process = move || { // TODO better way ? @@ -367,7 +371,7 @@ async fn process_sources(client: Arc, config: &Arc, Ok(rt) => { rt.block_on(async { let (input_stats, target_stats, mut res_errors) = - process_source(Arc::clone(&http_client), cfg, index, usr_trgts, event_manager).await; + process_source(Arc::clone(&http_client), cfg, index, usr_trgts, event_manager, playlist_state.as_ref()).await; shared_errors.lock().await.append(&mut res_errors); let process_stats = SourceStats::new(input_stats, target_stats); shared_stats.lock().await.push(process_stats); @@ -381,7 +385,8 @@ async fn process_sources(client: Arc, config: &Arc, handles.drain(..).for_each(|handle| { let _ = handle.join(); }); } } else { - let (input_stats, target_stats, mut res_errors) = process_source(Arc::clone(&client), cfg, index, usr_trgts, event_manager).await; + let (input_stats, target_stats, mut res_errors) = + process_source(Arc::clone(&client), cfg, index, usr_trgts, event_manager, playlist_state).await; shared_errors.lock().await.append(&mut res_errors); let process_stats = SourceStats::new(input_stats, target_stats); shared_stats.lock().await.push(process_stats); @@ -460,13 +465,16 @@ fn flatten_groups(playlistgroups: Vec) -> Vec { sort_order } +#[allow(clippy::too_many_arguments)] async fn process_playlist_for_target(app_config: &AppConfig, client: Arc, playlists: &mut [FetchedPlaylist<'_>], target: &ConfigTarget, stats: &mut HashMap, errors: &mut Vec, - event_manager: Option>) -> Result<(), Vec> { + event_manager: Option>, + playlist_state: Option<&Arc> +) -> Result<(), Vec> { let pipe = get_processing_pipe(target); debug_if_enabled!("Processing order is {}", &target.processing_order); @@ -527,7 +535,7 @@ async fn process_playlist_for_target(app_config: &AppConfig, if process_watch(&config, &client, target, &flat_new_playlist) { step.tick("group watches"); } - let result = persist_playlist(app_config, &mut flat_new_playlist, flatten_tvguide(&new_epg).as_ref(), target).await; + let result = persist_playlist(app_config, &mut flat_new_playlist, flatten_tvguide(&new_epg).as_ref(), target, playlist_state).await; step.stop("Persisting playlists"); result } @@ -582,10 +590,10 @@ fn process_watch(cfg: &Config, client: &Arc, target: &ConfigTar } } -pub async fn exec_processing(client: Arc, app_config: Arc, targets: Arc, event_manager: Option>) { +pub async fn exec_processing(client: Arc, app_config: Arc, targets: Arc, event_manager: Option>, playlist_state: Option>) { let start_time = Instant::now(); let event_manager_clone = event_manager.clone(); - let (stats, errors) = process_sources(Arc::clone(&client), &app_config, targets.clone(), event_manager_clone).await; + let (stats, errors) = process_sources(Arc::clone(&client), &app_config, targets.clone(), event_manager_clone, playlist_state.as_ref()).await; // log errors for err in &errors { error!("{}", err.message); diff --git a/backend/src/repository/m3u_playlist_iterator.rs b/backend/src/repository/m3u_playlist_iterator.rs index d527adeea..b0cdd1f3e 100644 --- a/backend/src/repository/m3u_playlist_iterator.rs +++ b/backend/src/repository/m3u_playlist_iterator.rs @@ -33,6 +33,9 @@ impl M3uPlaylistIterator { target: &ConfigTarget, user: &ProxyUserCredentials, ) -> Result { + + // TODO use playlist memory cache, but be aware of sorting ! + let m3u_output = target.get_m3u_output().ok_or_else(|| info_err!(format!("Unexpected failure, missing m3u target output for target {}", target.name)))?; let config = cfg.config.load(); let target_path = ensure_target_storage_path(&config, target.name.as_str())?; @@ -147,7 +150,6 @@ impl Iterator for M3uPlaylistIterator { pub struct M3uPlaylistM3uTextIterator { inner: M3uPlaylistIterator, started: bool, - } impl M3uPlaylistM3uTextIterator { diff --git a/backend/src/repository/m3u_repository.rs b/backend/src/repository/m3u_repository.rs index e0939308d..259bf1833 100644 --- a/backend/src/repository/m3u_repository.rs +++ b/backend/src/repository/m3u_repository.rs @@ -93,17 +93,12 @@ pub async fn m3u_get_item_for_stream_id(stream_id: u32, app_state: &AppState, ta return Err(str_to_io_error("id should start with 1")); } { - if let Some(playlist) = app_state.playlists.read().await.get(target.name.as_str()) { - return match playlist.m3u.as_ref() { - Some(m3u_playlist) => { - Ok(m3u_playlist.query(&stream_id) - .ok_or_else(|| str_to_io_error(&format!("Failed to read m3u item for id {stream_id}")))? - .clone()) - } - None => { - Err(str_to_io_error(&format!("Failed to read m3u item for id {stream_id}. It seems to be a xtream playlist"))) - } - }; + if let Some(playlist) = app_state.playlists.data.read().await.get(target.name.as_str()) { + if let Some(m3u_playlist) = playlist.m3u.as_ref() { + return Ok(m3u_playlist.query(&stream_id) + .ok_or_else(|| str_to_io_error(&format!("Failed to read m3u item for id {stream_id}")))? + .clone()) + } } let cfg: &AppConfig = &app_state.app_config; diff --git a/backend/src/repository/playlist_repository.rs b/backend/src/repository/playlist_repository.rs index 4649f4729..eb61e2678 100644 --- a/backend/src/repository/playlist_repository.rs +++ b/backend/src/repository/playlist_repository.rs @@ -1,24 +1,25 @@ -use shared::error::{info_err, TuliproxErrorKind}; -use shared::error::{TuliproxError}; -use crate::model::{AppConfig, ConfigTarget, TargetOutput}; -use shared::model::{PlaylistGroup, PlaylistItemType, XtreamCluster, XtreamPlaylistItem}; +use crate::api::model::{AppState, PlaylistM3uStorage, PlaylistStorage, PlaylistStorageState, PlaylistXtreamStorage}; use crate::model::Epg; +use crate::model::{AppConfig, ConfigTarget, TargetOutput}; +use crate::repository::bplustree::BPlusTree; use crate::repository::epg_repository::epg_write; -use crate::repository::strm_repository::write_strm_playlist; -use crate::repository::m3u_repository::m3u_write_playlist; +use crate::repository::indexed_document::IndexedDocumentIterator; +use crate::repository::m3u_repository::{m3u_get_file_paths, m3u_write_playlist}; use crate::repository::storage::{ensure_target_storage_path, get_target_id_mapping_file, get_target_storage_path}; +use crate::repository::strm_repository::write_strm_playlist; use crate::repository::target_id_mapping::{TargetIdMapping, VirtualIdRecord}; use crate::repository::xtream_repository::{xtream_get_file_paths, xtream_get_storage_path, xtream_write_playlist}; -use std::path::Path; -use shared::{create_tuliprox_error}; -use shared::utils::{is_dash_url, is_hls_url}; -use crate::api::model::{AppState, PlaylistStorage, PlaylistXtreamStorage}; -use crate::repository::bplustree::{BPlusTree}; -use crate::repository::indexed_document::{IndexedDocumentIterator}; use crate::utils; +use shared::error::{info_err, TuliproxErrorKind}; +use shared::error::TuliproxError; +use shared::model::{M3uPlaylistItem, PlaylistGroup, PlaylistItemType, XtreamCluster, XtreamPlaylistItem}; +use shared::utils::{is_dash_url, is_hls_url}; +use shared::create_tuliprox_error; +use std::path::Path; +use std::sync::Arc; pub async fn persist_playlist(app_config: &AppConfig, playlist: &mut [PlaylistGroup], epg: Option<&Epg>, - target: &ConfigTarget) -> Result<(), Vec> { + target: &ConfigTarget, playlist_state: Option<&Arc>) -> Result<(), Vec> { let mut errors = vec![]; let config = &app_config.config.load(); let target_path = match ensure_target_storage_path(config, &target.name) { @@ -60,12 +61,33 @@ pub async fn persist_playlist(app_config: &AppConfig, playlist: &mut [PlaylistGr TargetOutput::HdHomeRun(_hdhomerun_output) => Ok(()), }; - if let Err(err) = result { - errors.push(err); - } else if !playlist.is_empty() { - if let Err(err) = epg_write(config, target, &target_path, epg, output) { - errors.push(err); + match result { + Ok(()) => { + if !playlist.is_empty() { + if let Err(err) = epg_write(config, target, &target_path, epg, output) { + errors.push(err); + } + } + + if target.use_memory_cache { + if let Some(playlist_storage) = playlist_state { + match output { + TargetOutput::Xtream(_) => { + if let Ok(storage) = load_xtream_target_storage(app_config, target) { + playlist_storage.cache_playlist(&target.name, PlaylistStorage::XtreamPlaylist(Box::new(storage))).await; + } + }, + TargetOutput::M3u(_) => { + if let Ok(storage) = load_m3u_target_storage(app_config, target) { + playlist_storage.cache_playlist(&target.name, PlaylistStorage::M3uPlaylist(Box::new(storage))).await; + } + }, + _ => {} + } + } + } } + Err(err) => errors.push(err) } } @@ -85,28 +107,25 @@ pub async fn get_target_id_mapping(cfg: &AppConfig, target_path: &Path) -> (Targ fn load_target_id_mapping_as_tree(app_config: &AppConfig, target_path: &Path, target: &ConfigTarget) -> Result, TuliproxError> { - let target_id_mapping_file = get_target_id_mapping_file(&target_path); + let target_id_mapping_file = get_target_id_mapping_file(target_path); let _file_lock = app_config.file_locks.read_lock(&target_id_mapping_file); - Ok(BPlusTree::::load(&target_id_mapping_file).map_err(|err| + BPlusTree::::load(&target_id_mapping_file).map_err(|err| create_tuliprox_error!( TuliproxErrorKind::Info, "Could not find path for target {} err:{err}", &target.name - ))?) + )) } fn load_xtream_playlist_as_tree(app_config: &AppConfig, storage_path: &Path, cluster: XtreamCluster) -> BPlusTree { let (main_path, index_path) = xtream_get_file_paths(storage_path, cluster); let _file_lock = app_config.file_locks.read_lock(&main_path); let mut tree = BPlusTree::::new(); - match IndexedDocumentIterator::::new(&main_path, &index_path) { - Ok(reader) => { - for (doc, _has_next) in reader { - tree.insert(doc.virtual_id, doc); - } + if let Ok(reader) = IndexedDocumentIterator::::new(&main_path, &index_path) { + for (doc, _has_next) in reader { + tree.insert(doc.virtual_id, doc); } - Err(_) => {} - }; + } tree } @@ -136,25 +155,44 @@ fn load_xtream_target_storage(app_config: &AppConfig, target: &ConfigTarget) -> }) } -pub async fn load_playlists_into_memory_cache(app_state: &AppState) -> Result<(), TuliproxError>{ +fn load_m3u_target_storage(app_config: &AppConfig, target: &ConfigTarget) -> Result { + let config = app_config.config.load(); + let target_path = get_target_storage_path(&config, target.name.as_str()).ok_or_else(|| + create_tuliprox_error!( + TuliproxErrorKind::Info, + "Could not find path for target {}", &target.name + ))?; + + let (main_path, index_path) = m3u_get_file_paths(&target_path); + let _file_lock = app_config.file_locks.read_lock(&main_path); + let mut tree = BPlusTree::::new(); + if let Ok(reader) = IndexedDocumentIterator::::new(&main_path, &index_path) { + for (doc, _has_next) in reader { + tree.insert(doc.virtual_id, doc); + } + } + Ok(tree) +} + +pub async fn load_playlists_into_memory_cache(app_state: &AppState) -> Result<(), TuliproxError> { let app_config: &AppConfig = &app_state.app_config; - for sources in app_state.app_config.sources.load().sources.iter() { - for target in sources.targets.iter() { + for sources in &app_state.app_config.sources.load().sources { + for target in &sources.targets { if target.use_memory_cache { - for output in target.output.iter() { + for output in &target.output { match output { TargetOutput::Xtream(_) => { if let Ok(storage) = load_xtream_target_storage(app_config, target) { - app_state.cache_playlist(&target.name, PlaylistStorage::XtreamPlaylist(storage)).await; + app_state.cache_playlist(&target.name, PlaylistStorage::XtreamPlaylist(Box::new(storage))).await; } } TargetOutput::M3u(_) => { if let Ok(storage) = load_m3u_target_storage(app_config, target) { - app_state.cache_playlist(&target.name, PlaylistStorage::M3uPlaylist(storage)).await; + app_state.cache_playlist(&target.name, PlaylistStorage::M3uPlaylist(Box::new(storage))).await; } } _ => {} - }; + } }; } } diff --git a/backend/src/repository/xtream_playlist_iterator.rs b/backend/src/repository/xtream_playlist_iterator.rs index 47c4a6445..122a0c795 100644 --- a/backend/src/repository/xtream_playlist_iterator.rs +++ b/backend/src/repository/xtream_playlist_iterator.rs @@ -30,6 +30,9 @@ impl XtreamPlaylistIterator { category_id: Option, user: &ProxyUserCredentials, ) -> Result { + + // TODO use playlist memory cache and keep sorted + let xtream_output = target.get_xtream_output().ok_or_else(|| info_err!(format!("Unexpected: xtream output required for target {}", target.name)))?; let config = app_config.config.load(); if let Some(storage_path) = xtream_get_storage_path(&config, target.name.as_str()) { diff --git a/backend/src/repository/xtream_repository.rs b/backend/src/repository/xtream_repository.rs index 9c140abc4..694b6161f 100644 --- a/backend/src/repository/xtream_repository.rs +++ b/backend/src/repository/xtream_repository.rs @@ -348,11 +348,11 @@ async fn xtream_get_item_for_stream_id_from_memory( app_state: &AppState, target: &ConfigTarget, xtream_cluster: Option, -) -> Result<(XtreamPlaylistItem, VirtualIdRecord), Error> { - if let Some(playlist) = app_state.playlists.read().await.get(target.name.as_str()) { +) -> Result, Error> { + if let Some(playlist) = app_state.playlists.data.read().await.get(target.name.as_str()) { return match playlist.xtream.as_ref() { None => { - Err(str_to_io_error(&format!("Failed to read xtream item for id {virtual_id}. It seems to be a m3u playlist"))) + Ok(None) } Some(xtream_storage) => { let mapping = xtream_storage.id_mapping.query(&virtual_id).ok_or_else(|| str_to_io_error(&format!("Could not find mapping for target {} and id {}", target.name, virtual_id)))?.clone(); @@ -400,11 +400,12 @@ async fn xtream_get_item_for_stream_id_from_memory( } }; - result.map(|xpli| (xpli, mapping)) + result.map(|xpli| Some((xpli, mapping))) } }; } - Err(str_to_io_error(&format!("Failed to read xtream item for id {virtual_id}. No entry found."))) + //Err(str_to_io_error(&format!("Failed to read xtream item for id {virtual_id}. No entry found."))) + Ok(None) } pub async fn xtream_get_item_for_stream_id( @@ -414,10 +415,10 @@ pub async fn xtream_get_item_for_stream_id( xtream_cluster: Option, ) -> Result<(XtreamPlaylistItem, VirtualIdRecord), Error> { if target.use_memory_cache { - return xtream_get_item_for_stream_id_from_memory(virtual_id, - app_state, - target, - xtream_cluster).await; + if let Some((playlist_item, virtual_record)) = + xtream_get_item_for_stream_id_from_memory(virtual_id, app_state, target, xtream_cluster).await? { + return Ok((playlist_item, virtual_record)); + } } let app_config: &AppConfig = &app_state.app_config; diff --git a/backend/src/utils/network/xtream.rs b/backend/src/utils/network/xtream.rs index 11e46be71..cf8dae4d2 100644 --- a/backend/src/utils/network/xtream.rs +++ b/backend/src/utils/network/xtream.rs @@ -240,8 +240,6 @@ pub async fn get_xtream_playlist(cfg: &Config, client: Arc, inp } } } - // why we need a sort if there is no sort defined ? - //playlist_groups.sort_by(|a, b| a.title.partial_cmp(&b.title).unwrap_or(Ordering::Greater)); for (grp_id, plg) in (1_u32..).zip(playlist_groups.iter_mut()) { plg.id = grp_id;