diff --git a/backend/src/api/config_watch.rs b/backend/src/api/config_watch.rs index 4d88d2759..7d90a8c9d 100644 --- a/backend/src/api/config_watch.rs +++ b/backend/src/api/config_watch.rs @@ -43,8 +43,8 @@ impl ConfigFile { Ok(()) } - fn load_api_proxy(app_state: &Arc) -> Result<(), TuliproxError> { - match utils::read_api_proxy_config(&app_state.app_config, true) { + async fn load_api_proxy(app_state: &Arc) -> Result<(), TuliproxError> { + match utils::read_api_proxy_config(&app_state.app_config, true).await { Ok(Some(api_proxy)) => { app_state.app_config.set_api_proxy(api_proxy)?; let paths = > as Access>::load(&app_state.app_config.paths); @@ -101,7 +101,7 @@ impl ConfigFile { match self { ConfigFile::ApiProxy => { app_state.event_manager.send_event(EventMessage::ConfigChange(ConfigType::ApiProxy)); - ConfigFile::load_api_proxy(app_state) + ConfigFile::load_api_proxy(app_state).await } ConfigFile::Mapping => { app_state.event_manager.send_event(EventMessage::ConfigChange(ConfigType::Mapping)); diff --git a/backend/src/api/endpoints/v1_api.rs b/backend/src/api/endpoints/v1_api.rs index a5b80b0b0..086d0fd40 100644 --- a/backend/src/api/endpoints/v1_api.rs +++ b/backend/src/api/endpoints/v1_api.rs @@ -79,7 +79,7 @@ async fn geoip_update(axum::extract::State(app_state): axum::extract::State( None } -fn xtream_get_info_resource_url<'a>( +async fn xtream_get_info_resource_url<'a>( config: &'a AppConfig, pli: &'a XtreamPlaylistItem, target: &'a ConfigTarget, @@ -570,12 +570,12 @@ fn xtream_get_info_resource_url<'a>( config, target.name.as_str(), pli.get_virtual_id(), - ), + ).await, XtreamCluster::Series => xtream_repository::xtream_load_series_info( config, target.name.as_str(), pli.get_virtual_id(), - ), + ).await, XtreamCluster::Live => None, }; if let Some(content) = info_content { @@ -649,7 +649,7 @@ fn get_season_info_doc(doc: &Vec, season_id: u32) -> Option<&Value> { None } -fn xtream_get_season_resource_url<'a>( +async fn xtream_get_season_resource_url<'a>( config: &'a AppConfig, pli: &'a XtreamPlaylistItem, target: &'a ConfigTarget, @@ -660,7 +660,7 @@ fn xtream_get_season_resource_url<'a>( config, target.name.as_str(), pli.get_virtual_id(), - ), + ).await, XtreamCluster::Video | XtreamCluster::Live => None, }; if let Some(content) = info_content { @@ -733,14 +733,14 @@ async fn xtream_player_api_resource( &pli, &target, resource - )) + ).await) } else if resource.starts_with(crate::model::XC_SEASON_RESOURCE_PREFIX) { try_result_bad_request!(xtream_get_season_resource_url( &app_state.app_config, &pli, &target, resource - )) + ).await) } else { pli.get_field(resource) }; diff --git a/backend/src/api/main_api.rs b/backend/src/api/main_api.rs index 96a029f3d..2eed2389b 100644 --- a/backend/src/api/main_api.rs +++ b/backend/src/api/main_api.rs @@ -59,7 +59,7 @@ async fn healthcheck() -> impl axum::response::IntoResponse { axum::Json(create_healthcheck()) } -fn create_shared_data( +async fn create_shared_data( app_config: &Arc, forced_targets: &Arc, ) -> AppState { @@ -68,7 +68,7 @@ fn create_shared_data( let use_geoip = config.is_geoip_enabled(); let geoip = if use_geoip { let path = get_geoip_path(&config.working_dir); - let _file_lock = app_config.file_locks.read_lock(&path); + let _file_lock = app_config.file_locks.read_lock(&path).await; match GeoIp::load(&path) { Ok(db) => { info!("GeoIp db loaded"); @@ -249,7 +249,7 @@ pub async fn start_server( if web_ui_enabled { infos.push(format!("Web root: {}", web_dir_path.display())); } - let app_shared_data = create_shared_data(&app_config, &targets); + let app_shared_data = create_shared_data(&app_config, &targets).await; let app_state = Arc::new(app_shared_data); let shared_data = Arc::clone(&app_state); diff --git a/backend/src/api/model/app_state.rs b/backend/src/api/model/app_state.rs index c2bd74f10..7287862d7 100644 --- a/backend/src/api/model/app_state.rs +++ b/backend/src/api/model/app_state.rs @@ -282,7 +282,7 @@ impl AppState { if changes.geoip { let new_geoip = if use_geoip { let path = get_geoip_path(&working_dir); - let _file_lock = self.app_config.file_locks.read_lock(&path); + let _file_lock = self.app_config.file_locks.read_lock(&path).await; GeoIp::load(&path).ok().map(Arc::new) } else { None diff --git a/backend/src/main.rs b/backend/src/main.rs index efdf7c59c..13a79a69f 100644 --- a/backend/src/main.rs +++ b/backend/src/main.rs @@ -106,15 +106,15 @@ fn main() { if let Some(bts) = BUILD_TIMESTAMP.to_string().parse::>().ok().map(|datetime| datetime.format("%Y-%m-%d %H:%M:%S %Z").to_string()) { info!("Build time: {bts}"); } - - let app_config = utils::read_initial_app_config(&mut config_paths, true, true, args.server).unwrap_or_else(|err| exit!("{}", err)); - print_info(&app_config); - - let sources = > as Access>::load(&app_config.sources); - let targets = sources.validate_targets(args.target.as_ref()).unwrap_or_else(|err| exit!("{}", err)); - let rt = tokio::runtime::Runtime::new().unwrap(); let () = rt.block_on(async { + + let app_config = utils::read_initial_app_config(&mut config_paths, true, true, args.server).await.unwrap_or_else(|err| exit!("{}", err)); + print_info(&app_config); + + let sources = > as Access>::load(&app_config.sources); + let targets = sources.validate_targets(args.target.as_ref()).unwrap_or_else(|err| exit!("{}", err)); + if args.server { start_in_server_mode(Arc::new(app_config), Arc::new(targets)).await; } else { diff --git a/backend/src/model/config/api_proxy.rs b/backend/src/model/config/api_proxy.rs index 15fd97ce7..108d3e43a 100644 --- a/backend/src/model/config/api_proxy.rs +++ b/backend/src/model/config/api_proxy.rs @@ -100,14 +100,14 @@ impl ApiProxyConfig { // we have the option to store user in the config file or in the user_db // When we switch from one to other we need to migrate the existing data. /// # Panics - pub fn migrate_api_user(&mut self, cfg: &AppConfig, errors: &mut Vec) { + pub async fn migrate_api_user(&mut self, cfg: &AppConfig, errors: &mut Vec) { let paths = > as Access>::load(&cfg.paths); let api_proxy_file = paths.api_proxy_file_path.as_str(); if self.use_user_db { // we have user defined in config file. // we migrate them to the db and delete them from the config file if !&self.user.is_empty() { - if let Err(err) = merge_api_user(cfg, &self.user) { + if let Err(err) = merge_api_user(cfg, &self.user).await { errors.push(err.to_string()); } else { let config = > as Access>::load(&cfg.config); @@ -118,7 +118,7 @@ impl ApiProxyConfig { } } } - match load_api_user(cfg) { + match load_api_user(cfg).await { Ok(users) => { self.user = users; } @@ -132,7 +132,7 @@ impl ApiProxyConfig { if user_db_path.exists() { // we cant have user defined in db file. // we need to load them and save them into the config file - if let Ok(stored_users) = load_api_user(cfg) { + if let Ok(stored_users) = load_api_user(cfg).await { for stored_user in stored_users { if let Some(target_user) = self.user.iter_mut().find(|t| t.target == stored_user.target) { for stored_credential in &stored_user.credentials { @@ -151,7 +151,7 @@ impl ApiProxyConfig { if let Err(err) = save_api_proxy(api_proxy_file, backup_dir.as_ref(), &ApiProxyConfigDto::from(&*self)) { errors.push(format!("Error saving api proxy file: {err}")); } else { - backup_api_user_db_file(cfg, &user_db_path); + backup_api_user_db_file(cfg, &user_db_path).await; let _ = fs::remove_file(&user_db_path); } } diff --git a/backend/src/processing/processor/xtream.rs b/backend/src/processing/processor/xtream.rs index 4cdeb8ed6..f7591d133 100644 --- a/backend/src/processing/processor/xtream.rs +++ b/backend/src/processing/processor/xtream.rs @@ -108,7 +108,7 @@ where }; { - let file_lock = cfg.file_locks.read_lock(&file_path); + let file_lock = cfg.file_locks.read_lock(&file_path).await; if let Ok(info_records) = BPlusTree::::load(&file_path) { info_records.iter().for_each(|(provider_id, record)| { processed_info_ids.insert(*provider_id, extract_ts(record)); diff --git a/backend/src/processing/processor/xtream_series.rs b/backend/src/processing/processor/xtream_series.rs index 752d2ec55..2fe3cbcfe 100644 --- a/backend/src/processing/processor/xtream_series.rs +++ b/backend/src/processing/processor/xtream_series.rs @@ -142,7 +142,7 @@ async fn process_series_info( return result; }; - let _file_lock = app_config.file_locks.read_lock(&info_path); + let _file_lock = app_config.file_locks.read_lock(&info_path).await; // Contains the Series Info with episode listing let Ok(mut info_reader) = IndexedDocumentReader::::new(&info_path, &idx_path) else { return result; }; diff --git a/backend/src/processing/processor/xtream_vod.rs b/backend/src/processing/processor/xtream_vod.rs index df779515b..3122f9081 100644 --- a/backend/src/processing/processor/xtream_vod.rs +++ b/backend/src/processing/processor/xtream_vod.rs @@ -148,7 +148,7 @@ pub async fn playlist_resolve_vod(app_config: &AppConfig, client: Arc>(&content).ok().and_then(|info_doc| { info_doc.get("info").cloned().map(|info_content| { diff --git a/backend/src/repository/m3u_repository.rs b/backend/src/repository/m3u_repository.rs index 12a384e38..f494f60d6 100644 --- a/backend/src/repository/m3u_repository.rs +++ b/backend/src/repository/m3u_repository.rs @@ -62,7 +62,7 @@ pub async fn m3u_write_playlist(cfg: &AppConfig, target: &ConfigTarget, target_o persist_m3u_playlist_as_text(&cfg.config.load(), target, target_output, &m3u_playlist); { - let _file_lock = cfg.file_locks.write_lock(&m3u_path); + let _file_lock = cfg.file_locks.write_lock(&m3u_path).await; match IndexedDocumentWriter::new(m3u_path.clone(), idx_path) { Ok(mut writer) => { for m3u in m3u_playlist { @@ -105,7 +105,7 @@ pub async fn m3u_get_item_for_stream_id(stream_id: u32, app_state: &AppState, ta let cfg: &AppConfig = &app_state.app_config; let target_path = get_target_storage_path(&cfg.config.load(), target.name.as_str()).ok_or_else(|| str_to_io_error(&format!("Could not find path for target {}", &target.name)))?; let (m3u_path, idx_path) = m3u_get_file_paths(&target_path); - let _file_lock = cfg.file_locks.read_lock(&m3u_path); + let _file_lock = cfg.file_locks.read_lock(&m3u_path).await; IndexedDocumentDirectAccess::read_indexed_item::(&m3u_path, &idx_path, &stream_id) } } diff --git a/backend/src/repository/playlist_repository.rs b/backend/src/repository/playlist_repository.rs index 8b4f06a59..eaeb154d3 100644 --- a/backend/src/repository/playlist_repository.rs +++ b/backend/src/repository/playlist_repository.rs @@ -98,12 +98,12 @@ pub async fn persist_playlist(app_config: &AppConfig, playlist: &mut [PlaylistGr for output in &target.output { match output { TargetOutput::Xtream(_) => { - if let Ok(storage) = load_xtream_target_storage(app_config, target) { + if let Ok(storage) = load_xtream_target_storage(app_config, target).await { 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) { + if let Ok(storage) = load_m3u_target_storage(app_config, target).await { playlist_storage.cache_playlist(&target.name, PlaylistStorage::M3uPlaylist(Box::new(storage))).await; } }, @@ -123,9 +123,9 @@ 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> { +async 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 _file_lock = app_config.file_locks.read_lock(&target_id_mapping_file); + let _file_lock = app_config.file_locks.read_lock(&target_id_mapping_file).await; BPlusTree::::load(&target_id_mapping_file).map_err(|err| create_tuliprox_error!( @@ -134,9 +134,9 @@ fn load_target_id_mapping_as_tree(app_config: &AppConfig, target_path: &Path, ta )) } -fn load_xtream_playlist_as_tree(app_config: &AppConfig, storage_path: &Path, cluster: XtreamCluster) -> BPlusTree { +async 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 _file_lock = app_config.file_locks.read_lock(&main_path).await; let mut tree = BPlusTree::::new(); if let Ok(reader) = IndexedDocumentIterator::::new(&main_path, &index_path) { for (doc, _has_next) in reader { @@ -146,7 +146,7 @@ fn load_xtream_playlist_as_tree(app_config: &AppConfig, storage_path: &Path, clu tree } -fn load_xtream_target_storage(app_config: &AppConfig, target: &ConfigTarget) -> Result { +async fn load_xtream_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!( @@ -159,10 +159,10 @@ fn load_xtream_target_storage(app_config: &AppConfig, target: &ConfigTarget) -> TuliproxErrorKind::Info, "Could not find path for target {} xtream output", &target.name))?; - let target_id_mapping = load_target_id_mapping_as_tree(app_config, &target_path, target)?; - let live_storage = load_xtream_playlist_as_tree(app_config, &storage_path, XtreamCluster::Live); - let vod_storage = load_xtream_playlist_as_tree(app_config, &storage_path, XtreamCluster::Video); - let series_storage = load_xtream_playlist_as_tree(app_config, &storage_path, XtreamCluster::Series); + let target_id_mapping = load_target_id_mapping_as_tree(app_config, &target_path, target).await?; + let live_storage = load_xtream_playlist_as_tree(app_config, &storage_path, XtreamCluster::Live).await; + let vod_storage = load_xtream_playlist_as_tree(app_config, &storage_path, XtreamCluster::Video).await; + let series_storage = load_xtream_playlist_as_tree(app_config, &storage_path, XtreamCluster::Series).await; Ok(PlaylistXtreamStorage { id_mapping: target_id_mapping, @@ -172,7 +172,7 @@ fn load_xtream_target_storage(app_config: &AppConfig, target: &ConfigTarget) -> }) } -fn load_m3u_target_storage(app_config: &AppConfig, target: &ConfigTarget) -> Result { +async 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!( @@ -181,7 +181,7 @@ fn load_m3u_target_storage(app_config: &AppConfig, target: &ConfigTarget) -> Res ))?; let (main_path, index_path) = m3u_get_file_paths(&target_path); - let _file_lock = app_config.file_locks.read_lock(&main_path); + let _file_lock = app_config.file_locks.read_lock(&main_path).await; let mut tree = BPlusTree::::new(); if let Ok(reader) = IndexedDocumentIterator::::new(&main_path, &index_path) { for (doc, _has_next) in reader { @@ -207,12 +207,12 @@ pub async fn load_target_into_memory_cache(app_state: &AppState, target: &Arc { - if let Ok(storage) = load_xtream_target_storage(&app_state.app_config, target) { + if let Ok(storage) = load_xtream_target_storage(&app_state.app_config, target).await { app_state.cache_playlist(&target.name, PlaylistStorage::XtreamPlaylist(Box::new(storage))).await; } } TargetOutput::M3u(_) => { - if let Ok(storage) = load_m3u_target_storage(&app_state.app_config, target) { + if let Ok(storage) = load_m3u_target_storage(&app_state.app_config, target).await { app_state.cache_playlist(&target.name, PlaylistStorage::M3uPlaylist(Box::new(storage))).await; } } diff --git a/backend/src/repository/user_repository.rs b/backend/src/repository/user_repository.rs index 04022c2d6..20f1a52a7 100644 --- a/backend/src/repository/user_repository.rs +++ b/backend/src/repository/user_repository.rs @@ -120,23 +120,23 @@ fn add_target_user_to_user_tree(target_users: &[TargetUser], user_tree: &mut BPl } } -pub fn merge_api_user(cfg: &AppConfig, target_users: &[TargetUser]) -> Result { +pub async fn merge_api_user(cfg: &AppConfig, target_users: &[TargetUser]) -> Result { let path = get_api_user_db_path(cfg); - let lock = cfg.file_locks.read_lock(&path); + let lock = cfg.file_locks.read_lock(&path).await; let mut user_tree: BPlusTree = BPlusTree::load(&path).unwrap_or_else(|_| BPlusTree::new()); drop(lock); add_target_user_to_user_tree(target_users, &mut user_tree); - let _lock = cfg.file_locks.write_lock(&path); + let _lock = cfg.file_locks.write_lock(&path).await; user_tree.store(&path) } /// # Panics /// /// Will panic if `backup_dir` is not given -pub fn backup_api_user_db_file(cfg: &AppConfig, path: &Path) { +pub async fn backup_api_user_db_file(cfg: &AppConfig, path: &Path) { if let Some(backup_dir) = cfg.config.load().backup_dir.as_ref() { let backup_path = PathBuf::from(backup_dir).join(format!("{}_{}", storage_const::API_USER_DB_FILE, Local::now().format("%Y%m%d_%H%M%S"))); - let _lock = cfg.file_locks.read_lock(path); + let _lock = cfg.file_locks.read_lock(path).await; match std::fs::copy(path, &backup_path) { Ok(_) => {} Err(err) => { error!("Could not backup file {}:{}", &backup_path.to_str().unwrap_or("?"), err) } @@ -144,19 +144,19 @@ pub fn backup_api_user_db_file(cfg: &AppConfig, path: &Path) { } } -pub fn store_api_user(cfg: &AppConfig, target_users: &[TargetUser]) -> Result { +pub async fn store_api_user(cfg: &AppConfig, target_users: &[TargetUser]) -> Result { let mut user_tree = BPlusTree::::new(); add_target_user_to_user_tree(target_users, &mut user_tree); let path = get_api_user_db_path(cfg); - backup_api_user_db_file(cfg, &path); - let _lock = cfg.file_locks.write_lock(&path); + backup_api_user_db_file(cfg, &path).await; + let _lock = cfg.file_locks.write_lock(&path).await; user_tree.store(&path) } // TODO remove me if we get stable on user_db -pub fn load_api_user_deprecated(cfg: &AppConfig) -> Result, Error> { +pub async fn load_api_user_deprecated(cfg: &AppConfig) -> Result, Error> { let path = get_api_user_db_path(cfg); - let lock = cfg.file_locks.read_lock(&path); + let lock = cfg.file_locks.read_lock(&path).await; let user_tree = BPlusTree::::load(&path)?; drop(lock); let mut target_users: HashMap = HashMap::new(); @@ -180,10 +180,10 @@ pub fn load_api_user_deprecated(cfg: &AppConfig) -> Result, Erro } -pub fn load_api_user(cfg: &AppConfig) -> Result, Error> { +pub async fn load_api_user(cfg: &AppConfig) -> Result, Error> { let path = get_api_user_db_path(cfg); - let lock = cfg.file_locks.read_lock(&path); - let Ok(user_tree) = BPlusTree::::load(&path) else { return load_api_user_deprecated(cfg) }; + let lock = cfg.file_locks.read_lock(&path).await; + let Ok(user_tree) = BPlusTree::::load(&path) else { return load_api_user_deprecated(cfg).await }; drop(lock); let mut target_users: HashMap = HashMap::new(); for (_uname, stored_user) in &user_tree { diff --git a/backend/src/repository/xtream_repository.rs b/backend/src/repository/xtream_repository.rs index 9fafbdd8c..ad741a087 100644 --- a/backend/src/repository/xtream_repository.rs +++ b/backend/src/repository/xtream_repository.rs @@ -101,7 +101,7 @@ pub fn xtream_get_record_file_path(storage_path: &Path, item_type: PlaylistItemT _ => None, } } -fn write_playlists_to_file( +async fn write_playlists_to_file( cfg: &AppConfig, storage_path: &Path, collections: Vec<(XtreamCluster, &[&mut PlaylistItem])>, @@ -109,7 +109,7 @@ fn write_playlists_to_file( for (cluster, playlist) in collections { let (xtream_path, idx_path) = xtream_get_file_paths(storage_path, cluster); { - let _file_lock = cfg.file_locks.write_lock(&xtream_path); + let _file_lock = cfg.file_locks.write_lock(&xtream_path).await; match IndexedDocumentWriter::new(xtream_path.clone(), idx_path) { Ok(mut writer) => { for item in playlist { @@ -183,7 +183,7 @@ pub fn xtream_get_file_paths(storage_path: &Path, cluster: XtreamCluster) -> (Pa // xtream_get_file_paths_for_name(storage_path, storage_const::FILE_SERIES) // } -fn xtream_garbage_collect(config: &AppConfig, target_name: &str) -> std::io::Result<()> { +async fn xtream_garbage_collect(config: &AppConfig, target_name: &str) -> std::io::Result<()> { // Garbage collect series let storage_path = { let cfg = config.config.load(); @@ -194,7 +194,7 @@ fn xtream_garbage_collect(config: &AppConfig, target_name: &str) -> std::io::Res XtreamCluster::Series )); { - let _file_lock = config.file_locks.write_lock(&info_path); + let _file_lock = config.file_locks.write_lock(&info_path).await; IndexedDocumentGarbageCollector::::new(info_path.clone(), idx_path)?.garbage_collect()?; } Ok(()) @@ -273,9 +273,9 @@ pub async fn xtream_write_playlist( (XtreamCluster::Video, &vod_col), (XtreamCluster::Series, &series_col), ], - ) { + ).await { Ok(()) => { - if let Err(err) = xtream_garbage_collect(cfg, &target.name) { + if let Err(err) = xtream_garbage_collect(cfg, &target.name).await { if err.kind() != ErrorKind::NotFound { errors.push(format!("Garbage collection failed:{err}")); } @@ -311,7 +311,7 @@ pub fn xtream_get_collection_path( Err(str_to_io_error(&format!("Cant find collection: {target_name}/{collection_name}"))) } -fn xtream_read_item_for_stream_id( +async fn xtream_read_item_for_stream_id( cfg: &AppConfig, stream_id: u32, storage_path: &Path, @@ -319,19 +319,19 @@ fn xtream_read_item_for_stream_id( ) -> Result { let (xtream_path, idx_path) = xtream_get_file_paths(storage_path, cluster); { - let _file_lock = cfg.file_locks.read_lock(&xtream_path); + let _file_lock = cfg.file_locks.read_lock(&xtream_path).await; IndexedDocumentDirectAccess::read_indexed_item::(&xtream_path, &idx_path, &stream_id) } } -fn xtream_read_series_item_for_stream_id( +async fn xtream_read_series_item_for_stream_id( cfg: &AppConfig, stream_id: u32, storage_path: &Path, ) -> Result { let (xtream_path, idx_path) = xtream_get_file_paths(storage_path, XtreamCluster::Series); { - let _file_lock = cfg.file_locks.read_lock(&xtream_path); + let _file_lock = cfg.file_locks.read_lock(&xtream_path).await; IndexedDocumentDirectAccess::read_indexed_item::(&xtream_path, &idx_path, &stream_id) } } @@ -431,33 +431,33 @@ pub async fn xtream_get_item_for_stream_id( let storage_path = xtream_get_storage_path(&config, target.name.as_str()).ok_or_else(|| str_to_io_error(&format!("Could not find path for target {} xtream output", &target.name)))?; { 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); + let _file_lock = app_config.file_locks.read_lock(&target_id_mapping_file).await; let mut target_id_mapping = BPlusTreeQuery::::try_new(&target_id_mapping_file).map_err(|err| str_to_io_error(&format!("Could not load id mapping for target {} err:{err}", target.name)))?; let mapping = target_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)))?; let result = match mapping.item_type { PlaylistItemType::SeriesInfo => { - xtream_read_series_item_for_stream_id(app_config, virtual_id, &storage_path) + xtream_read_series_item_for_stream_id(app_config, virtual_id, &storage_path).await } PlaylistItemType::Series => { - if let Ok(mut item) = xtream_read_series_item_for_stream_id(app_config, mapping.parent_virtual_id, &storage_path) { + if let Ok(mut item) = xtream_read_series_item_for_stream_id(app_config, mapping.parent_virtual_id, &storage_path).await { item.provider_id = mapping.provider_id; item.item_type = PlaylistItemType::Series; Ok(item) } else { - xtream_read_item_for_stream_id(app_config, virtual_id, &storage_path, XtreamCluster::Series) + xtream_read_item_for_stream_id(app_config, virtual_id, &storage_path, XtreamCluster::Series).await } } PlaylistItemType::Catchup => { let cluster = try_cluster!(xtream_cluster, mapping.item_type, virtual_id)?; - let mut item = xtream_read_item_for_stream_id(app_config, mapping.parent_virtual_id, &storage_path, cluster)?; + let mut item = xtream_read_item_for_stream_id(app_config, mapping.parent_virtual_id, &storage_path, cluster).await?; item.provider_id = mapping.provider_id; item.item_type = PlaylistItemType::Catchup; Ok(item) } _ => { let cluster = try_cluster!(xtream_cluster, mapping.item_type, virtual_id)?; - xtream_read_item_for_stream_id(app_config, virtual_id, &storage_path, cluster) + xtream_read_item_for_stream_id(app_config, virtual_id, &storage_path, cluster).await } }; @@ -475,7 +475,7 @@ pub async fn xtream_load_rewrite_playlist( XtreamPlaylistJsonIterator::new(cluster, config, target, category_id, user).await } -pub fn xtream_write_series_info( +pub async fn xtream_write_series_info( app_config: &AppConfig, target_name: &str, series_info_id: u32, @@ -490,14 +490,14 @@ pub fn xtream_write_series_info( )); { - let _file_lock = app_config.file_locks.write_lock(&info_path); + let _file_lock = app_config.file_locks.write_lock(&info_path).await; let mut writer = IndexedDocumentWriter::new_append(info_path.clone(), idx_path)?; writer.write_doc(series_info_id, content).map_err(|_| str_to_io_error(&format!("failed to write xtream series info for target {target_name}")))?; writer.store()?; } { let target_id_mapping_file = get_target_id_mapping_file(&target_path); - let _file_lock = app_config.file_locks.write_lock(&target_id_mapping_file); + let _file_lock = app_config.file_locks.write_lock(&target_id_mapping_file).await; if let Ok(mut target_id_mapping) = BPlusTreeUpdate::::try_new(&target_id_mapping_file) { if let Some(record) = target_id_mapping.query(&series_info_id) { let new_record = record.copy_update_timestamp(); @@ -527,20 +527,20 @@ pub async fn xtream_write_vod_info( Ok(()) } -fn xtream_get_series_info_mapping( +async fn xtream_get_series_info_mapping( config: &AppConfig, target_name: &str, series_id: u32, ) -> Option { - xtream_get_info_mapping(config, target_name, series_id).filter(|id_record| !id_record.is_expired()) + xtream_get_info_mapping(config, target_name, series_id).await.filter(|id_record| !id_record.is_expired()) } -fn xtream_get_info_mapping(app_config: &AppConfig, target_name: &str, info_id: u32) -> Option { +async fn xtream_get_info_mapping(app_config: &AppConfig, target_name: &str, info_id: u32) -> Option { let config = app_config.config.load(); let target_path = get_target_storage_path(&config, target_name)?; 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); + let _file_lock = app_config.file_locks.read_lock(&target_id_mapping_file).await; BPlusTreeQuery::::try_new(&target_id_mapping_file).map_err(|err| { error!("Could not load id mapping for target {target_name}: {err}"); str_to_io_error(&format!("ID mapping load error for target {target_name}")) @@ -548,12 +548,12 @@ fn xtream_get_info_mapping(app_config: &AppConfig, target_name: &str, info_id: u } // Reads the series info entry if exists -pub fn xtream_load_series_info( +pub async fn xtream_load_series_info( app_config: &AppConfig, target_name: &str, series_id: u32, ) -> Option { - xtream_get_series_info_mapping(app_config, target_name, series_id)?; + xtream_get_series_info_mapping(app_config, target_name, series_id).await?; let config = app_config.config.load(); let storage_path = xtream_get_storage_path(&config, target_name)?; @@ -561,7 +561,7 @@ pub fn xtream_load_series_info( if info_path.exists() && idx_path.exists() { { - let _file_lock = app_config.file_locks.read_lock(&info_path); + let _file_lock = app_config.file_locks.read_lock(&info_path).await; return match IndexedDocumentDirectAccess::read_indexed_item::(&info_path, &idx_path, &series_id) { Ok(content) => Some(content), Err(err) => { @@ -573,24 +573,24 @@ pub fn xtream_load_series_info( } None } -fn xtream_get_vod_info_mapping( +async fn xtream_get_vod_info_mapping( config: &AppConfig, target_name: &str, vod_id: u32, ) -> Option { - xtream_get_info_mapping(config, target_name, vod_id) + xtream_get_info_mapping(config, target_name, vod_id).await //.filter(|id_record| !id_record.is_expired()) } // Reads the vod info entry if exists -pub fn xtream_load_vod_info( +pub async fn xtream_load_vod_info( config: &AppConfig, target_name: &str, vod_id: u32, ) -> Option { // Check if the entry exists; if not, we don't need to look further. - xtream_get_vod_info_mapping(config, target_name, vod_id).as_ref()?; + xtream_get_vod_info_mapping(config, target_name, vod_id).await.as_ref()?; // Entry exists, read db entry let target_storage_path = xtream_get_storage_path(&config.config.load(), target_name)?; @@ -598,7 +598,7 @@ pub fn xtream_load_vod_info( if info_path.exists() && idx_path.exists() { { - let _file_lock = config.file_locks.read_lock(&info_path); + let _file_lock = config.file_locks.read_lock(&info_path).await; return IndexedDocumentDirectAccess::read_indexed_item::(&info_path, &idx_path, &vod_id).ok(); } } @@ -809,11 +809,11 @@ pub async fn write_and_get_xtream_series_info

( let mut doc = serde_json::from_str::>(content).map_err(|_| str_to_io_error("Failed to parse JSON content"))?; let virtual_id = pli_series_info.get_virtual_id(); let app_config = &app_state.app_config; - xtream_write_series_info(app_config, target.name.as_str(), virtual_id, content).ok(); + xtream_write_series_info(app_config, target.name.as_str(), virtual_id, content).await.ok(); rewrite_xtream_series_info(app_state, target, xtream_output, pli_series_info, user, &mut doc).await } -pub fn xtream_get_input_info( +pub async fn xtream_get_input_info( cfg: &AppConfig, input: &ConfigInput, provider_id: u32, @@ -821,7 +821,7 @@ pub fn xtream_get_input_info( ) -> Option { if let Ok(Some((info_path, idx_path))) = get_input_storage_path(&input.name, &cfg.config.load().working_dir).map(|storage_path| xtream_get_info_file_paths(&storage_path, cluster)) { - let _file_lock = cfg.file_locks.read_lock(&info_path); + let _file_lock = cfg.file_locks.read_lock(&info_path).await; if let Ok(content) = IndexedDocumentDirectAccess::read_indexed_item::(&info_path, &idx_path, &provider_id) { return Some(content); } @@ -839,7 +839,7 @@ pub async fn xtream_update_input_info_file( match get_input_storage_path(&input.name, &config.working_dir).map(|storage_path| xtream_get_info_file_paths(&storage_path, cluster)) { Ok(Some((info_path, idx_path))) => { { - let _file_lock = cfg.file_locks.write_lock(&info_path); + let _file_lock = cfg.file_locks.write_lock(&info_path).await; let mut reader = file_reader(open_readonly_file(wal_path).map_err(|err| notify_err!(format!("Could not read {cluster} info {err}")))?); match IndexedDocumentWriter::::new_append(info_path.clone(), idx_path) { Ok(mut writer) => { @@ -885,7 +885,7 @@ pub async fn xtream_update_input_vod_record_from_wal_file( .and_then(|opt| opt.ok_or_else(|| notify_err!(format!("Error accessing storage path for input: {}", &input.name))))?; { - let _file_lock = cfg.file_locks.write_lock(&record_path); + let _file_lock = cfg.file_locks.write_lock(&record_path).await; let mut reader = file_reader(open_readonly_file(wal_path).map_err(|err| notify_err!(format!("Could not read vod wal info {err}")))?); let mut provider_id_bytes = [0u8; 4]; let mut tmdb_id_bytes = [0u8; 4]; @@ -953,7 +953,7 @@ pub async fn xtream_update_input_series_record_from_wal_file( .map_err(|err| notify_err!(format!("Error accessing storage path: {err}"))) .and_then(|opt| opt.ok_or_else(|| notify_err!(format!("Error accessing storage path for input: {}", &input.name))))?; { - let _file_lock = cfg.file_locks.write_lock(&record_path); + let _file_lock = cfg.file_locks.write_lock(&record_path).await; let mut reader = file_reader(open_readonly_file(wal_path).map_err(|err| notify_err!(format!("Could not read series wal info {err}")))?); let mut provider_id_bytes = [0u8; 4]; let mut ts_bytes = [0u8; 8]; @@ -988,7 +988,7 @@ pub async fn xtream_update_input_series_episodes_record_from_wal_file( .map_err(|err| notify_err!(format!("Error accessing storage path: {err}"))) .and_then(|opt| opt.ok_or_else(|| notify_err!(format!("Error accessing storage path for input: {}", &input.name))))?; { - let _file_lock = cfg.file_locks.write_lock(&record_path); + let _file_lock = cfg.file_locks.write_lock(&record_path).await; let mut reader = file_reader(open_readonly_file(wal_path).map_err(|err| notify_err!(format!("Could not read series episode wal info {err}")))?); let mut provider_id_bytes = [0u8; 4]; let mut len_bytes = [0u8; 4]; diff --git a/backend/src/utils/file/config_reader.rs b/backend/src/utils/file/config_reader.rs index d78595eaa..4d39e30e2 100644 --- a/backend/src/utils/file/config_reader.rs +++ b/backend/src/utils/file/config_reader.rs @@ -43,13 +43,13 @@ pub fn config_file_reader(file: File, resolve_env: bool) -> impl Read } } -pub fn read_api_proxy_config(config: &AppConfig, resolve_env: bool) -> Result, TuliproxError> { +pub async fn read_api_proxy_config(config: &AppConfig, resolve_env: bool) -> Result, TuliproxError> { let paths = > as Access>::load(&config.paths); let api_proxy_file_path = paths.api_proxy_file_path.as_str(); if let Some(api_proxy_dto) = read_api_proxy_file(api_proxy_file_path, resolve_env)? { let mut errors = vec![]; let mut api_proxy: ApiProxyConfig = ApiProxyConfig::from(&api_proxy_dto); - api_proxy.migrate_api_user(config, &mut errors); + api_proxy.migrate_api_user(config, &mut errors).await; if !errors.is_empty() { for error in errors { error!("{error}"); @@ -181,7 +181,7 @@ pub fn get_batch_aliases(input_type: InputType, url: &str) -> Result Result<(), TuliproxError> { +pub async fn prepare_users(app_config_dto: &mut AppConfigDto, app_config: &AppConfig) -> Result<(), TuliproxError> { let use_user_db = app_config_dto .api_proxy .as_ref() @@ -190,7 +190,7 @@ pub fn prepare_users(app_config_dto: &mut AppConfigDto, app_config: &AppConfig) if use_user_db { let user_db_path = get_api_user_db_path(app_config); if user_db_path.exists() { - match load_api_user(app_config) { + match load_api_user(app_config).await { Ok(stored_users) => if let Some(api_proxy) = app_config_dto.api_proxy.as_mut() { api_proxy.user.extend(stored_users.iter().map(TargetUserDto::from)); }, @@ -203,7 +203,7 @@ pub fn prepare_users(app_config_dto: &mut AppConfigDto, app_config: &AppConfig) Ok(()) } -pub fn read_initial_app_config(paths: &mut ConfigPaths, +pub async fn read_initial_app_config(paths: &mut ConfigPaths, resolve_env: bool, include_computed: bool, server_mode: bool) -> Result { @@ -250,7 +250,7 @@ pub fn read_initial_app_config(paths: &mut ConfigPaths, } if server_mode { - match read_api_proxy_config(&app_config, resolve_env) { + match read_api_proxy_config(&app_config, resolve_env).await { Ok(Some(api_proxy)) => app_config.set_api_proxy(api_proxy)?, Ok(None) => info!("Api-Proxy file: not used"), Err(err) => exit!("{err}"), @@ -279,13 +279,13 @@ pub fn read_api_proxy_file(api_proxy_file: &str, resolve_env: bool) -> Result Option { +pub async fn read_api_proxy(config: &AppConfig, resolve_env: bool) -> Option { let paths = > as Access>::load(&config.paths); match read_api_proxy_file(paths.api_proxy_file_path.as_str(), resolve_env) { Ok(Some(api_proxy_dto)) => { let mut errors = vec![]; let mut api_proxy: ApiProxyConfig = api_proxy_dto.into(); - api_proxy.migrate_api_user(config, &mut errors); + api_proxy.migrate_api_user(config, &mut errors).await; if !errors.is_empty() { for error in errors { error!("{error}"); diff --git a/backend/src/utils/network/xtream.rs b/backend/src/utils/network/xtream.rs index 3f0fdb243..723b434ce 100644 --- a/backend/src/utils/network/xtream.rs +++ b/backend/src/utils/network/xtream.rs @@ -70,7 +70,7 @@ where let app_config = &app_state.app_config; if cluster == XtreamCluster::Series { - if let Some(content) = xtream_repository::xtream_load_series_info(app_config, target.name.as_str(), pli.get_virtual_id()) { + if let Some(content) = xtream_repository::xtream_load_series_info(app_config, target.name.as_str(), pli.get_virtual_id()).await { // Deliver existing target content return rewrite_xtream_series_info_content(app_state, target, xtream_output, pli, user, &content).await; } @@ -78,20 +78,20 @@ where // Check if the content has been resolved if xtream_output.resolve_series { if let Some(provider_id) = pli.get_provider_id() { - if let Some(content) = xtream_get_input_info(app_config, input, provider_id, XtreamCluster::Series) { + if let Some(content) = xtream_get_input_info(app_config, input, provider_id, XtreamCluster::Series).await { return xtream_repository::write_and_get_xtream_series_info(app_state, target, xtream_output, pli, user, &content).await; } } } } else if cluster == XtreamCluster::Video { - if let Some(content) = xtream_repository::xtream_load_vod_info(app_config, target.name.as_str(), pli.get_virtual_id()) { + if let Some(content) = xtream_repository::xtream_load_vod_info(app_config, target.name.as_str(), pli.get_virtual_id()).await { // Deliver existing target content return rewrite_xtream_vod_info_content(app_config, target, xtream_output, pli, user, &content); } // Check if the content has been resolved if xtream_output.resolve_vod { if let Some(provider_id) = pli.get_provider_id() { - if let Some(content) = xtream_get_input_info(app_config, input, provider_id, XtreamCluster::Video) { + if let Some(content) = xtream_get_input_info(app_config, input, provider_id, XtreamCluster::Video).await { return xtream_repository::write_and_get_xtream_vod_info(app_config, target, xtream_output, pli, user, &content).await; } }