diff --git a/src/api/api_utils.rs b/src/api/api_utils.rs index 26958c5f5..c83968124 100644 --- a/src/api/api_utils.rs +++ b/src/api/api_utils.rs @@ -83,7 +83,7 @@ pub async fn stream_response(app_state: &AppState, stream_url: &str, req: &HttpR let send_bytes = matches!(item_type, PlaylistItemType::Video | PlaylistItemType::Series); if let Ok(url) = Url::parse(stream_url) { - let stream = buffered_stream::get_buffered_stream(&url, req, input, send_bytes); + let stream = buffered_stream::get_buffered_stream(&app_state.http_client, &url, req, input, send_bytes); return if share_stream { SharedStream::register(app_state, stream_url, stream).await; if let Some(broadcast_stream) = create_notify_stream(app_state, stream_url).await { @@ -128,7 +128,7 @@ pub fn get_headers_from_request(req: &HttpRequest) -> HashMap> { .collect() } -pub async fn resource_response(_app_state: &AppState, resource_url: &str, req: &HttpRequest, input: Option<&ConfigInput>) -> HttpResponse { +pub async fn resource_response(app_state: &AppState, resource_url: &str, req: &HttpRequest, input: Option<&ConfigInput>) -> HttpResponse { if resource_url.is_empty() { return HttpResponse::NoContent().finish(); } @@ -136,7 +136,7 @@ pub async fn resource_response(_app_state: &AppState, resource_url: &str, req: & debug_if_enabled!("Try to open resource {}", mask_sensitive_info(resource_url)); if let Ok(url) = Url::parse(resource_url) { - let client = request_utils::get_client_request(input.map(|i| &i.headers), &url, Some(&req_headers)); + let client = request_utils::get_client_request(&app_state.http_client, input.map(|i| &i.headers), &url, Some(&req_headers)); match client.send().await { Ok(response) => { let status = response.status(); diff --git a/src/api/download_api.rs b/src/api/download_api.rs index 195c6044b..68d5a620a 100644 --- a/src/api/download_api.rs +++ b/src/api/download_api.rs @@ -1,17 +1,17 @@ +use crate::api::model::app_state::AppState; +use crate::api::model::download::{DownloadQueue, FileDownload, FileDownloadRequest}; use crate::model::config::VideoDownloadConfig; use crate::utils::request_utils; use actix_web::{web, HttpResponse}; +use async_std::sync::RwLock; use futures::stream::TryStreamExt; use log::info; use serde_json::{json, Value}; use std::fs::File; use std::io::{ErrorKind, Write}; use std::ops::Deref; -use std::sync::{Arc}; -use async_std::sync::RwLock; +use std::sync::Arc; use std::{fs, io}; -use crate::api::model::app_state::AppState; -use crate::api::model::download::{DownloadQueue, FileDownload, FileDownloadRequest}; async fn download_file(active: Arc>>, client: &reqwest::Client) -> Result<(), String> { let file_download = { active.read().await.as_ref().unwrap().clone() }; diff --git a/src/api/main_api.rs b/src/api/main_api.rs index 2735d5589..a0a748a39 100644 --- a/src/api/main_api.rs +++ b/src/api/main_api.rs @@ -50,15 +50,16 @@ fn create_shared_data(cfg: &Arc) -> Data { finished: Arc::from(RwLock::new(Vec::new())), }), shared_streams: Arc::new(Mutex::new(HashMap::new())), + http_client: Arc::new(reqwest::Client::new()), }) } -fn exec_update_on_boot(cfg: &Arc, targets: &Arc) { +fn exec_update_on_boot(client: Arc, cfg: &Arc, targets: &Arc) { if cfg.update_on_boot { let cfg_clone = Arc::clone(cfg); let targets_clone = Arc::clone(targets); actix_rt::spawn( - async move { playlist_processor::exec_processing(cfg_clone, targets_clone).await } + async move { playlist_processor::exec_processing(client, cfg_clone, targets_clone).await } ); } } @@ -89,7 +90,7 @@ fn get_process_targets(cfg: &Arc, process_targets: &Arc, Arc::clone(process_targets) } -fn exec_scheduler(cfg: &Arc, targets: &Arc) { +fn exec_scheduler(client: &Arc, cfg: &Arc, targets: &Arc) { let schedules: Vec = if let Some(schedules) = &cfg.schedules { schedules.clone() } else { @@ -99,8 +100,9 @@ fn exec_scheduler(cfg: &Arc, targets: &Arc) { let expression = schedule.schedule.to_string(); let exec_targets = get_process_targets(cfg, targets, schedule.targets.as_ref()); let cfg_clone = Arc::clone(cfg); + let http_client = Arc::clone(client); actix_rt::spawn(async move { - start_scheduler(expression.as_str(), cfg_clone, exec_targets).await; + start_scheduler(http_client, expression.as_str(), cfg_clone, exec_targets).await; }); } } @@ -125,11 +127,12 @@ pub async fn start_server(cfg: Arc, targets: Arc) -> fut Err(err) => return Err(err) }; - exec_scheduler(&cfg, &targets); - exec_update_on_boot(&cfg, &targets); + let shared_data = create_shared_data(&cfg); + + exec_scheduler(&Arc::clone(&shared_data.http_client), &cfg, &targets); + exec_update_on_boot(Arc::clone(&shared_data.http_client), &cfg, &targets); let web_auth_enabled = is_web_auth_enabled(&cfg, web_ui_enabled, &web_dir_path); - let shared_data = create_shared_data(&cfg); // Web Server HttpServer::new(move || { App::new() diff --git a/src/api/model/app_state.rs b/src/api/model/app_state.rs index b184ef644..a419bbc58 100644 --- a/src/api/model/app_state.rs +++ b/src/api/model/app_state.rs @@ -9,4 +9,5 @@ pub struct AppState { pub config: Arc, pub downloads: Arc, pub shared_streams: Arc>>, + pub http_client: Arc, } diff --git a/src/api/model/buffered_stream.rs b/src/api/model/buffered_stream.rs index ed52f6270..4e3332021 100644 --- a/src/api/model/buffered_stream.rs +++ b/src/api/model/buffered_stream.rs @@ -1,11 +1,8 @@ use crate::api::api_utils::get_headers_from_request; use crate::model::config::ConfigInput; -use crate::utils::request_utils; -use crate::utils::request_utils::mask_sensitive_info; use actix_web::{HttpRequest, HttpResponse, HttpResponseBuilder}; use bytes::Bytes; use core::time::Duration; -use log::debug; use reqwest::header::RANGE; use reqwest::Error; use std::collections::HashMap; @@ -17,6 +14,7 @@ use tokio::sync::mpsc; use tokio::sync::mpsc::Receiver; use tokio_stream::Stream; use url::Url; +use crate::utils::request_utils::get_request_headers; const STREAM_QUEUE_SIZE: usize = 1024; // mpsc channel holding messages. const ERR_RETRY_TIMEOUT_SECS: u64 = 10; // If connect status is 4xx or 5xx, we wait until we allow next request from client @@ -70,20 +68,21 @@ impl Drop for BufferedReceiverStream { } } -pub fn get_buffered_stream(stream_url: &Url, req: &HttpRequest, input: Option<&ConfigInput>, range_send: bool) -> impl Stream> + Unpin + 'static { +pub fn get_buffered_stream(http_client: &Arc, stream_url: &Url, req: &HttpRequest, input: Option<&ConfigInput>, range_send: bool) -> impl Stream> + Unpin + 'static { let (tx, rx) = mpsc::channel::>(STREAM_QUEUE_SIZE); let req_headers = get_headers_from_request(req); let input_headers = input.map(|i| i.headers.clone()); let url = stream_url.clone(); let stop_signal = Arc::new(AtomicBool::new(false)); let stop_stream = Arc::clone(&stop_signal); - let req_client = request_utils::get_client_request(input_headers.as_ref(), &url, Some(&req_headers)); + let headers = get_request_headers(input_headers.as_ref(), Some(&req_headers)); + let base_client = Arc::clone(http_client); actix_rt::spawn(async move { - let masked_url = mask_sensitive_info(url.as_str()); + // let masked_url = mask_sensitive_info(url.as_str()); let req_bytes = get_request_bytes(&req_headers); let bytes_counter = if range_send { Some(AtomicUsize::new(req_bytes)) } else { None }; while !stop_signal.load(Ordering::Relaxed) { - let Some(mut client) = req_client.try_clone() else { break }; + let mut client = base_client.get(url.clone()).headers(headers.clone()); let bytes_to_request = bytes_counter.as_ref().map_or(0, |atomic| atomic.load(Ordering::Relaxed)); if bytes_to_request > 0 { // on reconnect send range header to avoid starting from beginning for vod diff --git a/src/api/model/shared_stream.rs b/src/api/model/shared_stream.rs index 80f7156fa..8aee72317 100644 --- a/src/api/model/shared_stream.rs +++ b/src/api/model/shared_stream.rs @@ -22,7 +22,7 @@ impl SharedStream { S: Stream> + Unpin + 'static, { // Create a broadcast channel for the shared stream - let (tx, _) = broadcast::channel(100); + let (tx, _) = broadcast::channel(10); let sender = Arc::new(tx); // Insert the shared stream into the shared state diff --git a/src/api/scheduler.rs b/src/api/scheduler.rs index 9809e6f43..d14aed530 100644 --- a/src/api/scheduler.rs +++ b/src/api/scheduler.rs @@ -24,7 +24,7 @@ fn datetime_to_instant(datetime: DateTime) -> Instant { Instant::now() + duration_until } -pub async fn start_scheduler(expression: &str, config: Arc, targets: Arc) -> ! { +pub async fn start_scheduler(client: Arc, expression: &str, config: Arc, targets: Arc) -> ! { match Schedule::from_str(expression) { Ok(schedule) => { let offset = *Local::now().offset(); @@ -32,7 +32,7 @@ pub async fn start_scheduler(expression: &str, config: Arc, targets: Arc let mut upcoming = schedule.upcoming(offset).take(1); if let Some(datetime) = upcoming.next() { actix_web::rt::time::sleep_until(actix_rt::time::Instant::from(datetime_to_instant(datetime))).await; - exec_processing(Arc::clone(&config), Arc::clone(&targets)).await; + exec_processing(Arc::clone(&client), Arc::clone(&config), Arc::clone(&targets)).await; } } } diff --git a/src/api/v1_api.rs b/src/api/v1_api.rs index edfd92f14..23a871041 100644 --- a/src/api/v1_api.rs +++ b/src/api/v1_api.rs @@ -103,7 +103,7 @@ async fn playlist_update( let process_targets = validate_targets(user_targets.as_ref(), &app_state.config.sources); match process_targets { Ok(valid_targets) => { - actix_rt::spawn(playlist_processor::exec_processing(Arc::clone(&app_state.config), Arc::new(valid_targets))); + actix_rt::spawn(playlist_processor::exec_processing(Arc::clone(&app_state.http_client), Arc::clone(&app_state.config), Arc::new(valid_targets))); HttpResponse::Ok().finish() } Err(err) => { @@ -128,13 +128,13 @@ fn create_config_input_for_url(url: &str) -> ConfigInput { } } -async fn get_playlist(cfg_input: Option<&ConfigInput>, cfg: &Config) -> HttpResponse { +async fn get_playlist(client: Arc, cfg_input: Option<&ConfigInput>, cfg: &Config) -> HttpResponse { match cfg_input { Some(input) => { let (result, errors) = match input.input_type { - InputType::M3u => download::get_m3u_playlist(cfg, input, &cfg.working_dir).await, - InputType::Xtream => download::get_xtream_playlist(input, &cfg.working_dir).await, + InputType::M3u => download::get_m3u_playlist(client, cfg, input, &cfg.working_dir).await, + InputType::Xtream => download::get_xtream_playlist(client, input, &cfg.working_dir).await, }; if result.is_empty() { let error_strings: Vec = errors.iter().map(std::string::ToString::to_string).collect(); @@ -152,11 +152,11 @@ async fn playlist( app_state: web::Data, ) -> HttpResponse { if let Some(input_id) = req.input_id { - get_playlist(app_state.config.get_input_by_id(input_id), &app_state.config).await + get_playlist(Arc::clone(&app_state.http_client), app_state.config.get_input_by_id(input_id), &app_state.config).await } else { let url = req.url.as_deref().unwrap_or(""); let input = create_config_input_for_url(url); - get_playlist(Some(&input), &app_state.config).await + get_playlist(Arc::clone(&app_state.http_client), Some(&input), &app_state.config).await } } diff --git a/src/api/xtream_api.rs b/src/api/xtream_api.rs index 929939ac6..df289c0ed 100644 --- a/src/api/xtream_api.rs +++ b/src/api/xtream_api.rs @@ -1,5 +1,6 @@ // https://github.com/tellytv/go.xtream-codes/blob/master/structs.go +use crate::Arc; use std::collections::HashMap; use std::fmt::{Display, Formatter}; use std::io::{Error, ErrorKind}; @@ -407,7 +408,7 @@ async fn xtream_get_stream_info_response(app_state: &AppState, user: &ProxyUserC // Redirect is only possible for live streams, vod and series info needs to be modified if user.proxy == ProxyType::Redirect && cluster == XtreamCluster::Live { return HttpResponse::Found().insert_header(("Location", info_url)).finish(); - } else if let Ok(content) = download::get_xtream_stream_info(&app_state.config, user, input, target, &pli, info_url.as_str(), cluster).await { + } else if let Ok(content) = download::get_xtream_stream_info(Arc::clone(&app_state.http_client), &app_state.config, user, input, target, &pli, info_url.as_str(), cluster).await { return HttpResponse::Ok().content_type(mime::APPLICATION_JSON).body(content); } } @@ -440,7 +441,7 @@ async fn xtream_get_short_epg(app_state: &AppState, user: &ProxyUserCredentials, return HttpResponse::Found().insert_header(("Location", info_url)).finish(); } - return match request_utils::download_text_content(input, info_url.as_str(), None).await { + return match request_utils::download_text_content(Arc::clone(&app_state.http_client), input, info_url.as_str(), None).await { Ok(content) => HttpResponse::Ok().content_type(mime::APPLICATION_JSON).body(content), Err(err) => { error!("Failed to download epg {}", mask_sensitive_info(err.to_string().as_str())); @@ -481,7 +482,7 @@ async fn xtream_get_catchup_response(app_state: &AppState, target: &ConfigTarget let pli = try_result_bad_request!(xtream_repository::xtream_get_item_for_stream_id(virtual_id, &app_state.config, target, Some(XtreamCluster::Live)).await); let input = try_option_bad_request!(app_state.config.get_input_by_id(pli.input_id)); let info_url = try_option_bad_request!(download::get_xtream_player_api_action_url(input, ACTION_GET_CATCHUP_TABLE).map(|action_url| format!("{action_url}&{TAG_STREAM_ID}={}&start={start}&end={end}", pli.provider_id))); - let content = try_result_bad_request!(download::get_xtream_stream_info_content(info_url.as_str(), input).await); + let content = try_result_bad_request!(download::get_xtream_stream_info_content(Arc::clone(&app_state.http_client), info_url.as_str(), input).await); let mut doc: Map = try_result_bad_request!(serde_json::from_str(&content)); let epg_listings = try_option_bad_request!(doc.get_mut(TAG_EPG_LISTINGS).and_then(Value::as_array_mut)); let target_path = try_option_bad_request!(get_target_storage_path(&app_state.config, target.name.as_str())); diff --git a/src/main.rs b/src/main.rs index 095fb667c..b1359ed49 100644 --- a/src/main.rs +++ b/src/main.rs @@ -1,21 +1,21 @@ #![allow(clippy::module_name_repetitions)] +extern crate core; extern crate env_logger; extern crate pest; #[macro_use] extern crate pest_derive; -extern crate core; +use actix_rt::System; use std::fs::File; use std::path::PathBuf; use std::sync::Arc; -use actix_rt::System; +use crate::auth::password::generate_password; use clap::Parser; use env_logger::Builder; use log::{error, info, LevelFilter}; -use crate::auth::password::generate_password; -use crate::model::config::{Config, HealthcheckConfig, ProcessTargets, validate_targets}; +use crate::model::config::{validate_targets, Config, HealthcheckConfig, ProcessTargets}; use crate::model::healthcheck::Healthcheck; use crate::processing::playlist_processor; use crate::utils::{config_reader, file_utils}; @@ -70,7 +70,8 @@ struct Args { #[arg(short = None, long = "genpwd", default_value_t = false, default_missing_value = "true")] genpwd: bool, - #[arg(short = None, long = "healthcheck", default_value_t = false, default_missing_value = "true")] + #[arg(short = None, long = "healthcheck", default_value_t = false, default_missing_value = "true" + )] healthcheck: bool, } @@ -100,7 +101,7 @@ fn main() { let sources_file: String = args.source_file.unwrap_or_else(|| file_utils::get_default_sources_file_path(&config_path)); let mut cfg = config_reader::read_config(config_path.as_str(), config_file.as_str(), sources_file.as_str()).unwrap_or_else(|err| exit!("{}", err)); - if args.genpwd { + if args.genpwd { match generate_password() { Ok(pwd) => println!("{pwd}"), Err(err) => error!("{err}") @@ -143,30 +144,30 @@ fn create_directories(cfg: &Config) { cfg.video.as_ref().and_then(|v| v.download.as_ref()).and_then(|d| d.directory.clone()) ]; - let mut paths : Vec= paths_strings.iter() + let mut paths: Vec = paths_strings.iter() .filter_map(|opt| opt.as_ref()) // Get rid of the `Option` .map(PathBuf::from).collect(); let mut temp_path = PathBuf::from(&cfg.working_dir); temp_path.push("tmp"); paths.push(temp_path); - // Iterate over the paths, filter out `None` values, and process the `Some(path)` values. + // Iterate over the paths, filter out `None` values, and process the `Some(path)` values. for path in &paths { - if !path.exists() { - // Create the directory tree if it doesn't exist - let path_value = path.to_str().unwrap_or("?"); - if let Err(e) = std::fs::create_dir_all(path) { - error!("Failed to create directory {path_value}: {e}"); - } else { - info!("Created directory: {path_value}"); - } + if !path.exists() { + // Create the directory tree if it doesn't exist + let path_value = path.to_str().unwrap_or("?"); + if let Err(e) = std::fs::create_dir_all(path) { + error!("Failed to create directory {path_value}: {e}"); + } else { + info!("Created directory: {path_value}"); } } - + } } fn start_in_cli_mode(cfg: Arc, targets: Arc) { - System::new().block_on(async { playlist_processor::exec_processing(cfg, targets).await }); + let client = Arc::new(reqwest::Client::new()); + System::new().block_on(async { playlist_processor::exec_processing(client, cfg, targets).await }); } fn start_in_server_mode(cfg: Arc, targets: Arc) { diff --git a/src/processing/playlist_processor.rs b/src/processing/playlist_processor.rs index 711090d1e..613f923ce 100644 --- a/src/processing/playlist_processor.rs +++ b/src/processing/playlist_processor.rs @@ -291,7 +291,7 @@ fn is_target_enabled(target: &ConfigTarget, user_targets: &ProcessTargets) -> bo (!user_targets.enabled && target.enabled) || (user_targets.enabled && user_targets.has_target(target.id)) } -async fn process_source(cfg: Arc, source_idx: usize, user_targets: Arc) -> (Vec, Vec, Vec) { +async fn process_source(client: Arc, cfg: Arc, source_idx: usize, user_targets: Arc) -> (Vec, Vec, Vec) { let source = cfg.sources.get(source_idx).unwrap(); let mut errors = vec![]; let mut input_stats = HashMap::::new(); @@ -304,11 +304,11 @@ async fn process_source(cfg: Arc, source_idx: usize, user_targets: Arc

download::get_m3u_playlist(&cfg, input, &cfg.working_dir).await, - InputType::Xtream => download::get_xtream_playlist(input, &cfg.working_dir).await, + InputType::M3u => download::get_m3u_playlist(Arc::clone(&client), &cfg, input, &cfg.working_dir).await, + InputType::Xtream => download::get_xtream_playlist(Arc::clone(&client), input, &cfg.working_dir).await, }; let (tvguide, mut tvguide_errors) = if error_list.is_empty() { - download::get_xmltv(&cfg, input, &cfg.working_dir).await + download::get_xmltv(Arc::clone(&client), &cfg, input, &cfg.working_dir).await } else { (None, vec![]) }; @@ -344,7 +344,7 @@ async fn process_source(cfg: Arc, source_idx: usize, user_targets: Arc

()); for target in &source.targets { if is_target_enabled(target, &user_targets) { - match process_playlist_for_target(&mut source_playlists, target, &cfg, &mut input_stats, &mut errors).await { + match process_playlist_for_target(Arc::clone(&client), &mut source_playlists, target, &cfg, &mut input_stats, &mut errors).await { Ok(()) => { target_stats.push(TargetStats::success(&target.name)); } @@ -376,7 +376,7 @@ fn create_input_stat(group_count: usize, channel_count: usize, error_count: usiz } } -async fn process_sources(config: Arc, user_targets: Arc) -> (Vec, Vec) { +async fn process_sources(client: Arc, config: Arc, user_targets: Arc) -> (Vec, Vec) { let mut handle_list = vec![]; let thread_num = config.threads; let process_parallel = thread_num > 1 && config.sources.len() > 1; @@ -398,10 +398,11 @@ async fn process_sources(config: Arc, user_targets: Arc) let cfg = config.clone(); let usr_trgts = user_targets.clone(); if process_parallel { + let http_client = Arc::clone(&client); let handles = &mut handle_list; let process = move || { System::new().block_on(async { - let (input_stats, target_stats, mut res_errors) = process_source(cfg, index, usr_trgts).await; + let (input_stats, target_stats, mut res_errors) = process_source(Arc::clone(&http_client), cfg, index, usr_trgts).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); @@ -412,7 +413,7 @@ async fn process_sources(config: Arc, user_targets: Arc) handles.drain(..).for_each(|handle| { let _ = handle.join(); }); } } else { - let (input_stats, target_stats, mut res_errors) = process_source(cfg, index, usr_trgts).await; + let (input_stats, target_stats, mut res_errors) = process_source(Arc::clone(&client), cfg, index, usr_trgts).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); @@ -485,7 +486,8 @@ fn flatten_groups(playlistgroups: Vec) -> Vec { sort_order } -async fn process_playlist_for_target(playlists: &mut [FetchedPlaylist<'_>], +async fn process_playlist_for_target(client: Arc, + playlists: &mut [FetchedPlaylist<'_>], target: &ConfigTarget, cfg: &Config, stats: &mut HashMap, @@ -497,8 +499,8 @@ async fn process_playlist_for_target(playlists: &mut [FetchedPlaylist<'_>], let mut processed_fetched_playlists: Vec = vec![]; for provider_fpl in playlists.iter_mut() { let mut processed_fpl = execute_pipe(target, &pipe, provider_fpl, &mut duplicates); - playlist_resolve_series(cfg, target, errors, &pipe, provider_fpl, &mut processed_fpl).await; - playlist_resolve_vod(cfg, target, errors, &processed_fpl).await; + playlist_resolve_series(Arc::clone(&client), cfg, target, errors, &pipe, provider_fpl, &mut processed_fpl).await; + playlist_resolve_vod(Arc::clone(&client), cfg, target, errors, &processed_fpl).await; // stats let input_stats = stats.get_mut(&processed_fpl.input.id); if let Some(stat) = input_stats { @@ -560,9 +562,9 @@ fn process_watch(target: &ConfigTarget, cfg: &Config, new_playlist: &Vec, targets: Arc) { +pub async fn exec_processing(client: Arc, cfg: Arc, targets: Arc) { let start_time = Instant::now(); - let (stats, errors) = process_sources(cfg.clone(), targets.clone()).await; + let (stats, errors) = process_sources(client, cfg.clone(), targets.clone()).await; // log errors for err in &errors { error!("{}", err.message); diff --git a/src/processing/xtream_processor.rs b/src/processing/xtream_processor.rs index 617de6e17..73852e5b5 100644 --- a/src/processing/xtream_processor.rs +++ b/src/processing/xtream_processor.rs @@ -9,6 +9,7 @@ use std::collections::HashMap; use std::fs::File; use std::io::{BufWriter, Error, ErrorKind, Write}; use std::path::PathBuf; +use std::sync::Arc; const FILE_SERIES_INFO: &str = "xtream_series_info"; const FILE_VOD_INFO: &str = "xtream_vod_info"; @@ -53,11 +54,11 @@ macro_rules! create_resolve_options_function_for_xtream_target { }; } -pub(in crate::processing) async fn playlist_resolve_download_playlist_item(pli: &PlaylistItem, input: &ConfigInput, errors: &mut Vec, resolve_delay: u16, cluster: XtreamCluster) -> Option { +pub(in crate::processing) async fn playlist_resolve_download_playlist_item(client: Arc, pli: &PlaylistItem, input: &ConfigInput, errors: &mut Vec, resolve_delay: u16, cluster: XtreamCluster) -> Option { let mut result = None; let provider_id = pli.get_provider_id()?; if let Some(info_url) = download::get_xtream_player_api_info_url(input, cluster, provider_id) { - result = match download::get_xtream_stream_info_content(&info_url, input).await { + result = match download::get_xtream_stream_info_content(client, &info_url, input).await { Ok(content) => Some(content), Err(err) => { errors.push(info_err!(format!("{err}"))); diff --git a/src/processing/xtream_processor_series.rs b/src/processing/xtream_processor_series.rs index 1b2423ecb..df80abf3b 100644 --- a/src/processing/xtream_processor_series.rs +++ b/src/processing/xtream_processor_series.rs @@ -11,6 +11,7 @@ use crate::{create_resolve_options_function_for_xtream_target, handle_error, han use std::collections::HashMap; use std::fs::File; use std::io::{BufWriter, Write}; +use std::sync::Arc; use std::time::Instant; use log::{info, log_enabled, Level}; use crate::model::xtream::{XtreamSeriesEpisode, XtreamSeriesInfoEpisode}; @@ -52,7 +53,7 @@ fn should_update_series_info(pli: &PlaylistItem, processed_provider_ids: &HashMa should_update_info(pli, processed_provider_ids, TAG_SERIES_INFO_LAST_MODIFIED) } -async fn playlist_resolve_series_info(cfg: &Config, errors: &mut Vec, +async fn playlist_resolve_series_info(client: Arc, cfg: &Config, errors: &mut Vec, fpl: &mut FetchedPlaylist<'_>, resolve_delay: u16) -> bool { let mut processed_info_ids = read_processed_series_info_ids(cfg, errors, fpl).await; // we cant write to the indexed-document directly because of the write lock and time-consuming operation. @@ -78,7 +79,7 @@ async fn playlist_resolve_series_info(cfg: &Config, errors: &mut Vec, cfg: &Config, target: &ConfigTarget, errors: &mut Vec, pipe: &ProcessingPipe, provider_fpl: &mut FetchedPlaylist<'_>, @@ -217,7 +218,7 @@ pub async fn playlist_resolve_series(cfg: &Config, target: &ConfigTarget, let (resolve_series, resolve_delay) = get_resolve_series_options(target, processed_fpl); if !resolve_series { return; } - if !playlist_resolve_series_info(cfg, errors, processed_fpl, resolve_delay).await { return; } + if !playlist_resolve_series_info(client, cfg, errors, processed_fpl, resolve_delay).await { return; } let series_playlist = process_series_info(cfg, provider_fpl, errors).await; if series_playlist.is_empty() { return; } // original content saved into original list diff --git a/src/processing/xtream_processor_vod.rs b/src/processing/xtream_processor_vod.rs index 981fc7f1d..eb54d5174 100644 --- a/src/processing/xtream_processor_vod.rs +++ b/src/processing/xtream_processor_vod.rs @@ -9,6 +9,7 @@ use serde_json::{Map, Value}; use std::collections::HashMap; use std::fs::File; use std::io::{BufWriter, Write}; +use std::sync::Arc; use std::time::Instant; use log::{info, log_enabled, Level}; @@ -63,7 +64,7 @@ fn should_update_vod_info(pli: &PlaylistItem, processed_provider_ids: &HashMap, fpl: &FetchedPlaylist<'_>) { +pub async fn playlist_resolve_vod(client: Arc, cfg: &Config, target: &ConfigTarget, errors: &mut Vec, fpl: &FetchedPlaylist<'_>) { let (resolve_movies, resolve_delay) = get_resolve_vod_options(target, fpl); if !resolve_movies { return; } @@ -92,7 +93,7 @@ pub async fn playlist_resolve_vod(cfg: &Config, target: &ConfigTarget, errors: & for pli in vod_info_iter { let (should_update, _provider_id, _ts) = should_update_vod_info(pli, &processed_info_ids); if should_update { - if let Some(content) = playlist_resolve_download_playlist_item(pli, fpl.input, errors, resolve_delay, XtreamCluster::Video).await { + if let Some(content) = playlist_resolve_download_playlist_item(Arc::clone(&client), pli, fpl.input, errors, resolve_delay, XtreamCluster::Video).await { if let Some((provider_id, info_record)) = extract_info_record_from_vod_info(&content) { let ts = info_record.ts; handle_error_and_return!(write_info_content_to_wal_file(&mut content_writer, provider_id, &content), diff --git a/src/utils/download.rs b/src/utils/download.rs index b3ddc1a28..3d56469a6 100644 --- a/src/utils/download.rs +++ b/src/utils/download.rs @@ -1,3 +1,4 @@ +use crate::Arc; use std::borrow::Cow; use crate::m3u_filter_error::M3uFilterError; use crate::model::config::{Config, ConfigInput, ConfigTarget}; @@ -30,10 +31,10 @@ fn prepare_file_path(persist: Option<&str>, working_dir: &str, action: &str) -> } } -pub async fn get_m3u_playlist(cfg: &Config, input: &ConfigInput, working_dir: &str) -> (Vec, Vec) { +pub async fn get_m3u_playlist(client: Arc, cfg: &Config, input: &ConfigInput, working_dir: &str) -> (Vec, Vec) { let url = input.url.clone(); let persist_file_path = prepare_file_path(input.persist.as_deref(), working_dir, ""); - match request_utils::get_input_text_content(input, working_dir, &url, persist_file_path).await { + match request_utils::get_input_text_content(client, input, working_dir, &url, persist_file_path).await { Ok(text) => { (m3u_parser::parse_m3u(cfg, input, text.lines()), vec![]) } @@ -64,12 +65,19 @@ pub fn get_xtream_player_api_info_url(input: &ConfigInput, cluster: XtreamCluste } -pub async fn get_xtream_stream_info_content(info_url: &str, input: &ConfigInput) -> Result { - request_utils::download_text_content(input, info_url, None).await +pub async fn get_xtream_stream_info_content(client: Arc, info_url: &str, input: &ConfigInput) -> Result { + request_utils::download_text_content(client, input, info_url, None).await } -pub async fn get_xtream_stream_info

(config: &Config, user: &ProxyUserCredentials, input: &ConfigInput, target: &ConfigTarget, - pli: &P, info_url: &str, cluster: XtreamCluster) -> Result +#[allow(clippy::too_many_arguments)] +pub async fn get_xtream_stream_info

(client: Arc, + config: &Config, + user: &ProxyUserCredentials, + input: &ConfigInput, + target: &ConfigTarget, + pli: &P, + info_url: &str, + cluster: XtreamCluster) -> Result where P: PlaylistEntry, { @@ -104,7 +112,7 @@ where } } - if let Ok(content) = get_xtream_stream_info_content(info_url, input).await { + if let Ok(content) = get_xtream_stream_info_content(client, info_url, input).await { return match cluster { XtreamCluster::Live => Ok(content), XtreamCluster::Video => xtream_repository::write_and_get_xtream_vod_info(config, target, pli, user, &content).await, @@ -141,7 +149,7 @@ const ACTIONS: [(XtreamCluster, &str, &str); 3] = [ (XtreamCluster::Video, "get_vod_categories", "get_vod_streams"), (XtreamCluster::Series, "get_series_categories", "get_series")]; -pub async fn get_xtream_playlist(input: &ConfigInput, working_dir: &str) -> (Vec, Vec) { +pub async fn get_xtream_playlist(client: Arc, input: &ConfigInput, working_dir: &str) -> (Vec, Vec) { let mut playlist_groups: Vec = Vec::with_capacity(128); let username = input.username.as_ref().map_or("", |v| v); let password = input.password.as_ref().map_or("", |v| v); @@ -158,8 +166,8 @@ pub async fn get_xtream_playlist(input: &ConfigInput, working_dir: &str) -> (Vec let stream_file_path = prepare_file_path(input.persist.as_deref(), working_dir, format!("{stream}_").as_str()); match futures::join!( - request_utils::get_input_json_content(input, category_url.as_str(), category_file_path), - request_utils::get_input_json_content(input, stream_url.as_str(), stream_file_path) + request_utils::get_input_json_content(Arc::clone(&client), input, category_url.as_str(), category_file_path), + request_utils::get_input_json_content(Arc::clone(&client), input, stream_url.as_str(), stream_file_path) ) { (Ok(category_content), Ok(stream_content)) => { match xtream_parser::parse_xtream(input, @@ -189,7 +197,7 @@ pub async fn get_xtream_playlist(input: &ConfigInput, working_dir: &str) -> (Vec (playlist_groups, errors) } -pub async fn get_xmltv(_cfg: &Config, input: &ConfigInput, working_dir: &str) -> (Option, Vec) { +pub async fn get_xmltv(client: Arc, _cfg: &Config, input: &ConfigInput, working_dir: &str) -> (Option, Vec) { match &input.epg_url { None => (None, vec![]), Some(url) => { @@ -197,7 +205,7 @@ pub async fn get_xmltv(_cfg: &Config, input: &ConfigInput, working_dir: &str) -> let persist_file_path = prepare_file_path(input.persist.as_deref(), working_dir, "") .map(|path| file_utils::add_prefix_to_filename(&path, "epg_", Some("xml"))); - match request_utils::get_input_text_content_as_file(input, working_dir, url, persist_file_path).await { + match request_utils::get_input_text_content_as_file(client, input, working_dir, url, persist_file_path).await { Ok(file) => { (Some(TVGuide { file }), vec![]) } diff --git a/src/utils/request_utils.rs b/src/utils/request_utils.rs index 49b36b248..416f10622 100644 --- a/src/utils/request_utils.rs +++ b/src/utils/request_utils.rs @@ -1,3 +1,4 @@ +use crate::Arc; use std::collections::{HashMap, HashSet}; use std::fs; use std::fs::File; @@ -27,10 +28,10 @@ pub const fn bytes_to_megabytes(bytes: u64) -> u64 { bytes / 1_048_576 } -pub async fn get_input_text_content_as_file(input: &ConfigInput, working_dir: &str, url_str: &str, persist_filepath: Option) -> Result { +pub async fn get_input_text_content_as_file(client: Arc, input: &ConfigInput, working_dir: &str, url_str: &str, persist_filepath: Option) -> Result { debug_if_enabled!("getting input text content working_dir: {}, url: {}", working_dir, mask_sensitive_info(url_str)); if url_str.parse::().is_ok() { - match download_text_content_as_file(input, url_str, working_dir, persist_filepath).await { + match download_text_content_as_file(client, input, url_str, working_dir, persist_filepath).await { Ok(content) => Ok(content), Err(e) => { error!("cant download input url: {} => {}", mask_sensitive_info(url_str), mask_sensitive_info(e.to_string().as_str())); @@ -73,11 +74,11 @@ pub async fn get_input_text_content_as_file(input: &ConfigInput, working_dir: &s } -pub async fn get_input_text_content(input: &ConfigInput, working_dir: &str, url_str: &str, persist_filepath: Option) -> Result { +pub async fn get_input_text_content(client: Arc, input: &ConfigInput, working_dir: &str, url_str: &str, persist_filepath: Option) -> Result { debug_if_enabled!("getting input text content working_dir: {}, url: {}", working_dir, mask_sensitive_info(url_str)); if url_str.parse::().is_ok() { - match download_text_content(input, url_str, persist_filepath).await { + match download_text_content(client, input, url_str, persist_filepath).await { Ok(content) => Ok(content), Err(e) => { error!("cant download input url: {} => {}", mask_sensitive_info(url_str), mask_sensitive_info(e.to_string().as_str())); @@ -119,8 +120,8 @@ pub async fn get_input_text_content(input: &ConfigInput, working_dir: &str, url_ } } -pub fn get_client_request(headers: Option<&HashMap>, url: &Url, custom_headers: Option<&HashMap>>) -> reqwest::RequestBuilder { - let request = reqwest::Client::new().get(url.clone()); +pub fn get_client_request(client: &Arc, headers: Option<&HashMap>, url: &Url, custom_headers: Option<&HashMap>>) -> reqwest::RequestBuilder { + let request = client.get(url.clone()); let headers = get_request_headers(headers, custom_headers); request.headers(headers) } @@ -176,9 +177,9 @@ fn get_local_file_content(file_path: &PathBuf) -> Result { } -async fn get_remote_content_as_file(input: &ConfigInput, url: &Url, file_path: &Path) -> Result { +async fn get_remote_content_as_file(client: Arc, input: &ConfigInput, url: &Url, file_path: &Path) -> Result { let start_time = Instant::now(); - let request = get_client_request(Some(&input.headers), url, None); + let request = get_client_request(&client, Some(&input.headers), url, None); match request.send().await { Ok(response) => { if response.status().is_success() { @@ -209,9 +210,9 @@ async fn get_remote_content_as_file(input: &ConfigInput, url: &Url, file_path: & } } -async fn get_remote_content(input: &ConfigInput, url: &Url) -> Result { +async fn get_remote_content(client: Arc, input: &ConfigInput, url: &Url) -> Result { let start_time = Instant::now(); - let request = get_client_request(Some(&input.headers), url, None); + let request = get_client_request(&client, Some(&input.headers), url, None); match request.send().await { Ok(response) => { let is_success = response.status().is_success(); @@ -272,7 +273,7 @@ async fn get_remote_content(input: &ConfigInput, url: &Url) -> Result) -> Result { +pub async fn download_text_content_as_file(client: Arc, input: &ConfigInput, url_str: &str, working_dir: &str, persist_filepath: Option) -> Result { if let Ok(url) = url_str.parse::() { if url.scheme() == "file" { url.to_file_path().map_or_else(|()| Err(Error::new(ErrorKind::Unsupported, format!("Unknown file {}", mask_sensitive_info(url_str)))), |file_path| if file_path.exists() { @@ -288,7 +289,7 @@ pub async fn download_text_content_as_file(input: &ConfigInput, url_str: &str, w Err(err) => Err(err) }, Ok); match file_path { - Ok(persist_path) => get_remote_content_as_file(input, &url, &persist_path).await, + Ok(persist_path) => get_remote_content_as_file(client, input, &url, &persist_path).await, Err(err) => Err(err) } } @@ -298,12 +299,12 @@ pub async fn download_text_content_as_file(input: &ConfigInput, url_str: &str, w } -pub async fn download_text_content(input: &ConfigInput, url_str: &str, persist_filepath: Option) -> Result { +pub async fn download_text_content(client: Arc, input: &ConfigInput, url_str: &str, persist_filepath: Option) -> Result { if let Ok(url) = url_str.parse::() { let result = if url.scheme() == "file" { url.to_file_path().map_or_else(|()| Err(Error::new(ErrorKind::Other, format!("Unknown file {}", mask_sensitive_info(url_str)))), |file_path| get_local_file_content(&file_path)) } else { - get_remote_content(input, &url).await + get_remote_content(client, input, &url).await }; match result { Ok(content) => { @@ -319,9 +320,9 @@ pub async fn download_text_content(input: &ConfigInput, url_str: &str, persist_f } } -async fn download_json_content(input: &ConfigInput, url: &str, persist_filepath: Option) -> Result { +async fn download_json_content(client: Arc, input: &ConfigInput, url: &str, persist_filepath: Option) -> Result { debug_if_enabled!("downloading json content from {}", mask_sensitive_info(url)); - match download_text_content(input, url, persist_filepath).await { + match download_text_content(client, input, url, persist_filepath).await { Ok(content) => { match serde_json::from_str::(&content) { Ok(value) => Ok(value), @@ -332,8 +333,8 @@ async fn download_json_content(input: &ConfigInput, url: &str, persist_filepath: } } -pub async fn get_input_json_content(input: &ConfigInput, url: &str, persist_filepath: Option) -> Result { - match download_json_content(input, url, persist_filepath).await { +pub async fn get_input_json_content(client: Arc, input: &ConfigInput, url: &str, persist_filepath: Option) -> Result { + match download_json_content(client, input, url, persist_filepath).await { Ok(content) => Ok(content), Err(e) => create_m3u_filter_error_result!(M3uFilterErrorKind::Notify, "cant download input url: {} => {}", mask_sensitive_info(url), mask_sensitive_info(e.to_string().as_str())) }