From 55191d6c84eeba9bf448352c22940ddf74119316 Mon Sep 17 00:00:00 2001 From: euzu Date: Thu, 11 Dec 2025 09:53:49 +0100 Subject: [PATCH] vod resolve refactored Arc is now Client --- .../src/api/endpoints/api_playlist_utils.rs | 6 +- backend/src/api/endpoints/hls_api.rs | 2 +- backend/src/api/endpoints/v1_api.rs | 2 +- backend/src/api/endpoints/v1_api_config.rs | 2 +- backend/src/api/endpoints/v1_api_playlist.rs | 11 +- backend/src/api/endpoints/xtream_api.rs | 6 +- backend/src/api/main_api.rs | 11 +- .../model/streams/provider_stream_factory.rs | 14 +-- backend/src/api/scheduler.rs | 9 +- backend/src/main.rs | 2 +- backend/src/messaging.rs | 13 +-- backend/src/model/xmltv.rs | 3 +- backend/src/processing/playlist_watch.rs | 5 +- backend/src/processing/processor/playlist.rs | 43 ++++--- backend/src/processing/processor/trakt.rs | 7 +- backend/src/processing/processor/xtream.rs | 3 +- .../src/processing/processor/xtream_series.rs | 7 +- .../src/processing/processor/xtream_vod.rs | 109 +++++++++--------- backend/src/utils/network/epg.rs | 9 +- backend/src/utils/network/ip_checker.rs | 8 +- backend/src/utils/network/m3u.rs | 2 +- backend/src/utils/network/request.rs | 29 +++-- backend/src/utils/network/xtream.rs | 20 ++-- backend/src/utils/telegram.rs | 3 +- backend/src/utils/trakt/client.rs | 5 +- 25 files changed, 163 insertions(+), 168 deletions(-) diff --git a/backend/src/api/endpoints/api_playlist_utils.rs b/backend/src/api/endpoints/api_playlist_utils.rs index 7658e5433..dabc8a398 100644 --- a/backend/src/api/endpoints/api_playlist_utils.rs +++ b/backend/src/api/endpoints/api_playlist_utils.rs @@ -141,13 +141,13 @@ pub(in crate::api::endpoints) async fn get_playlist_for_target(cfg_target: Optio (axum::http::StatusCode::BAD_REQUEST, axum::Json(json!({"error": "Invalid Arguments"}))).into_response() } -pub(in crate::api::endpoints) async fn get_playlist(client: Arc, cfg_input: Option<&Arc>, cfg: &Arc, accept: Option<&String>) -> impl IntoResponse + Send { +pub(in crate::api::endpoints) async fn get_playlist(client: &reqwest::Client, cfg_input: Option<&Arc>, cfg: &Arc, accept: Option<&String>) -> impl IntoResponse + Send { match cfg_input { Some(input) => { let (result, errors) = match input.input_type { - InputType::M3u | InputType::M3uBatch => m3u::get_m3u_playlist(&client, cfg, input, &cfg.working_dir).await, - InputType::Xtream | InputType::XtreamBatch => xtream::get_xtream_playlist(cfg, &client, input, &cfg.working_dir).await, + InputType::M3u | InputType::M3uBatch => m3u::get_m3u_playlist(client, cfg, input, &cfg.working_dir).await, + InputType::Xtream | InputType::XtreamBatch => xtream::get_xtream_playlist(cfg, client, input, &cfg.working_dir).await, }; if result.is_empty() { let error_strings: Vec = errors.iter().map(std::string::ToString::to_string).collect(); diff --git a/backend/src/api/endpoints/hls_api.rs b/backend/src/api/endpoints/hls_api.rs index c54e2148c..57f0fb53b 100644 --- a/backend/src/api/endpoints/hls_api.rs +++ b/backend/src/api/endpoints/hls_api.rs @@ -119,7 +119,7 @@ pub(in crate::api) async fn handle_hls_stream_request( let headers = request::get_request_headers(None, Some(&forwarded), disabled_headers.as_ref()); let input_source = InputSource::from(input).with_url(request_url); match request::download_text_content( - Arc::clone(&app_state.http_client.load()), + &app_state.http_client.load(), disabled_headers.as_ref(), &input_source, Some(&headers), diff --git a/backend/src/api/endpoints/v1_api.rs b/backend/src/api/endpoints/v1_api.rs index 523022c34..39ba81507 100644 --- a/backend/src/api/endpoints/v1_api.rs +++ b/backend/src/api/endpoints/v1_api.rs @@ -93,7 +93,7 @@ async fn geoip_update(axum::extract::State(app_state): axum::extract::State { let reader = Cursor::new(content); let mut geoip = GeoIp::new(); diff --git a/backend/src/api/endpoints/v1_api_config.rs b/backend/src/api/endpoints/v1_api_config.rs index 3bfb32e52..8a4a6c64f 100644 --- a/backend/src/api/endpoints/v1_api_config.rs +++ b/backend/src/api/endpoints/v1_api_config.rs @@ -123,7 +123,7 @@ async fn config_batch_content( .reverse_proxy .as_ref() .and_then(|r| r.disabled_header.clone()); - return match download_text_content(Arc::clone(&app_state.http_client.load()), disabled_headers.as_ref(), &input_source, None, None).await { + return match download_text_content(&app_state.http_client.load(), disabled_headers.as_ref(), &input_source, None, None).await { Ok((content, _path)) => { // Return CSV with explicit content-type try_unwrap_body!(axum::response::Response::builder() diff --git a/backend/src/api/endpoints/v1_api_playlist.rs b/backend/src/api/endpoints/v1_api_playlist.rs index 8ebbff432..db8f19ba7 100644 --- a/backend/src/api/endpoints/v1_api_playlist.rs +++ b/backend/src/api/endpoints/v1_api_playlist.rs @@ -60,14 +60,14 @@ async fn playlist_update( let process_targets = app_state.app_config.sources.load().validate_targets(user_targets.as_ref()); match process_targets { Ok(valid_targets) => { - let http_client = Arc::clone(&app_state.http_client.load()); + let http_client = app_state.http_client.load().as_ref().clone(); let app_config = Arc::clone(&app_state.app_config); let event_manager = Arc::clone(&app_state.event_manager); let playlist_state = Arc::clone(&app_state.playlists); let valid_targets = Arc::new(valid_targets); tokio::spawn({ async move { - playlist::exec_processing(http_client, app_config, valid_targets, Some(event_manager), Some(playlist_state)).await; + playlist::exec_processing(&http_client, app_config, valid_targets, Some(event_manager), Some(playlist_state)).await; } }); axum::http::StatusCode::ACCEPTED.into_response() @@ -86,18 +86,19 @@ async fn playlist_content( axum::extract::Json(playlist_req): axum::extract::Json, ) -> impl IntoResponse + Send { let config = app_state.app_config.config.load(); + let client = app_state.http_client.load(); match playlist_req { PlaylistRequest::Target(target_id) => { get_playlist_for_target(app_state.app_config.get_target_by_id(target_id).as_deref(), &app_state.app_config, accept.as_ref()).await.into_response() } PlaylistRequest::Input(input_id) => { - get_playlist(Arc::clone(&app_state.http_client.load()), app_state.app_config.get_input_by_id(input_id).as_ref(), &config, accept.as_ref()).await.into_response() + get_playlist(client.as_ref(), app_state.app_config.get_input_by_id(input_id).as_ref(), &config, accept.as_ref()).await.into_response() } PlaylistRequest::CustomXtream(xtream) => { match Url::parse(&xtream.url) { Ok(parsed) if parsed.scheme() == "http" || parsed.scheme() == "https" => { let input = Arc::new(create_config_input_for_xtream(&xtream.username, &xtream.password, &xtream.url)); - get_playlist(Arc::clone(&app_state.http_client.load()), Some(&input), &config, accept.as_ref()).await.into_response() + get_playlist(client.as_ref(), Some(&input), &config, accept.as_ref()).await.into_response() } _ => { (axum::http::StatusCode::BAD_REQUEST, axum::Json(json!({"error": "Invalid url scheme; only http/https are allowed"}))).into_response() @@ -108,7 +109,7 @@ async fn playlist_content( match Url::parse(&m3u.url) { Ok(parsed) if parsed.scheme() == "http" || parsed.scheme() == "https" => { let input = Arc::new(create_config_input_for_m3u(&m3u.url)); - get_playlist(Arc::clone(&app_state.http_client.load()), Some(&input), &config, accept.as_ref()).await.into_response() + get_playlist(client.as_ref(), Some(&input), &config, accept.as_ref()).await.into_response() } _ => { (axum::http::StatusCode::BAD_REQUEST, axum::Json(json!({"error": "Invalid url scheme; only http/https are allowed"}))).into_response() diff --git a/backend/src/api/endpoints/xtream_api.rs b/backend/src/api/endpoints/xtream_api.rs index eea853b46..897f55d07 100644 --- a/backend/src/api/endpoints/xtream_api.rs +++ b/backend/src/api/endpoints/xtream_api.rs @@ -981,7 +981,7 @@ async fn xtream_get_stream_info_response( if user.proxy == ProxyType::Redirect && cluster == XtreamCluster::Live { return redirect(&info_url).into_response(); } else if let Ok(content) = xtream::get_xtream_stream_info( - Arc::clone(&app_state.http_client.load()), + &app_state.http_client.load(), app_state, user, &input, @@ -1090,7 +1090,7 @@ async fn xtream_get_short_epg( // TODO serve epg from own db let input_source = InputSource::from(&*input).with_url(info_url); return match request::download_text_content( - Arc::clone(&app_state.http_client.load()), + &app_state.http_client.load(), None, &input_source, None, @@ -1243,7 +1243,7 @@ async fn xtream_get_catchup_response( let input_source = InputSource::from(&*input).with_url(info_url); let content = try_result_bad_request!( xtream::get_xtream_stream_info_content( - Arc::clone(&app_state.http_client.load()), + &app_state.http_client.load(), &input_source ) .await diff --git a/backend/src/api/main_api.rs b/backend/src/api/main_api.rs index 67e535ace..6d2948bbd 100644 --- a/backend/src/api/main_api.rs +++ b/backend/src/api/main_api.rs @@ -114,7 +114,7 @@ async fn create_shared_data( } fn exec_update_on_boot( - client: Arc, + client: &reqwest::Client, app_state: &Arc, targets: &Arc, ) { @@ -127,8 +127,9 @@ fn exec_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); + let client = client.clone(); tokio::spawn(async move { - playlist::exec_processing(client, app_state_clone, targets_clone, None, Some(playlist_state)).await; + playlist::exec_processing(&client, app_state_clone, targets_clone, None, Some(playlist_state)).await; }); } } @@ -272,15 +273,17 @@ pub async fn start_server( exec_system_usage(&app_state); + let client = shared_data.http_client.load(); + exec_scheduler( - &Arc::clone(&shared_data.http_client.load()), + client.as_ref(), &app_state, &targets, &cancel_token_scheduler, ); exec_update_on_boot( - Arc::clone(&shared_data.http_client.load()), + client.as_ref(), &app_state, &targets, ); diff --git a/backend/src/api/model/streams/provider_stream_factory.rs b/backend/src/api/model/streams/provider_stream_factory.rs index cf30108c5..d001649c4 100644 --- a/backend/src/api/model/streams/provider_stream_factory.rs +++ b/backend/src/api/model/streams/provider_stream_factory.rs @@ -193,7 +193,7 @@ fn get_request_range_start_bytes(req_headers: &HashMap>) -> Opti // } fn prepare_client( - request_client: &Arc, + request_client: &reqwest::Client, stream_options: &ProviderStreamFactoryOptions, ) -> (reqwest::RequestBuilder, bool) { let url = stream_options.get_url(); @@ -263,7 +263,7 @@ fn prepare_client( async fn provider_stream_request( app_state: &Arc, - request_client: &Arc, + request_client: &reqwest::Client, stream_options: &ProviderStreamFactoryOptions, ) -> Result, StatusCode> { let (client, _partial_content) = prepare_client(request_client, stream_options); @@ -365,7 +365,7 @@ async fn handle_channel_unavailable_stream(app_state: &Arc, async fn get_provider_stream( app_state: &Arc, - client: &Arc, + client: &reqwest::Client, stream_options: &ProviderStreamFactoryOptions, ) -> Result, StatusCode> { let url = stream_options.get_url(); @@ -419,7 +419,7 @@ async fn get_provider_stream( #[allow(clippy::too_many_lines)] pub async fn create_provider_stream( app_state: &Arc, - client: &Arc, + client: &reqwest::Client, stream_options: ProviderStreamFactoryOptions, ) -> Option { let client_stream_factory = |stream, reconnect_flag, range_cnt| { @@ -461,9 +461,9 @@ pub async fn create_provider_stream( let continue_streaming_signal = continue_client_signal.clone(); let stream_options_provider = stream_options.clone(); let app_state_clone = Arc::clone(app_state); - let client = Arc::clone(client); + let client = client.clone(); let unfold: BoxedProviderStream = stream::unfold((), move |()| { - let client = Arc::clone(&client); + let client = client.clone(); let stream_opts = stream_options_provider.clone(); let continue_streaming = continue_streaming_signal.clone(); let app_state_clone = Arc::clone(&app_state_clone); @@ -572,7 +572,7 @@ pub async fn create_provider_stream( // let input = None; // // let options = BufferStreamOptions::new(PlaylistItemType::Live, true, true, 0, false); -// let value = create_provider_stream(&cfg, Arc::clone(&client), &url, &req, input, options); +// let value = create_provider_stream(&cfg, &client, &url, &req, input, options); // let mut values = value.await; // 'outer: while let Some((ref mut stream, info)) = values.as_mut() { // if info.is_some() { diff --git a/backend/src/api/scheduler.rs b/backend/src/api/scheduler.rs index 0d8fa3768..7b511d092 100644 --- a/backend/src/api/scheduler.rs +++ b/backend/src/api/scheduler.rs @@ -26,7 +26,7 @@ pub fn datetime_to_instant(datetime: DateTime) -> Instant { Instant::now() + duration_until } -pub fn exec_scheduler(client: &Arc, app_state: &Arc, targets: &Arc, +pub fn exec_scheduler(client: &reqwest::Client, app_state: &Arc, targets: &Arc, cancel: &CancellationToken) { let cfg = &app_state.app_config; let config = cfg.config.load(); @@ -39,7 +39,7 @@ pub fn exec_scheduler(client: &Arc, app_state: &Arc, let expression = schedule.schedule.clone(); let exec_targets = get_process_targets(cfg, targets, schedule.targets.as_ref()); let app_state_clone = Arc::clone(app_state); - let http_client = Arc::clone(client); + let http_client = client.clone(); let cancel_token = cancel.clone(); tokio::spawn(async move { start_scheduler(http_client, expression.as_str(), app_state_clone, exec_targets, cancel_token).await; @@ -47,7 +47,7 @@ pub fn exec_scheduler(client: &Arc, app_state: &Arc, } } -async fn start_scheduler(client: Arc, expression: &str, app_state: Arc, +async fn start_scheduler(client: reqwest::Client, expression: &str, app_state: Arc, targets: Arc, cancel: CancellationToken) { match Schedule::from_str(expression) { Ok(schedule) => { @@ -55,12 +55,13 @@ async fn start_scheduler(client: Arc, expression: &str, app_sta loop { let mut upcoming = schedule.upcoming(offset).take(1); if let Some(datetime) = upcoming.next() { + let client = client.clone(); tokio::select! { () = 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); let playlist_state = app_state.playlists.clone(); - exec_processing(Arc::clone(&client), app_config, Arc::clone(&targets), Some(event_manager), Some(playlist_state)).await; + exec_processing(&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 706f29763..c186780aa 100644 --- a/backend/src/main.rs +++ b/backend/src/main.rs @@ -164,7 +164,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, None).await; + playlist::exec_processing(&client, cfg, targets, None, None).await; } async fn start_in_server_mode(cfg: Arc, targets: Arc) { diff --git a/backend/src/messaging.rs b/backend/src/messaging.rs index 8f7d6d610..54d6cd44b 100644 --- a/backend/src/messaging.rs +++ b/backend/src/messaging.rs @@ -5,13 +5,12 @@ use log::{debug, error}; use reqwest::header; use shared::model::MsgKind; use shared::utils::json_str_to_markdown; -use std::sync::Arc; fn is_enabled(kind: MsgKind, cfg: &MessagingConfig) -> bool { cfg.notify_on.contains(&kind) } -async fn send_http_post_request(client: &Arc, msg: &str, messaging: &MessagingConfig) { +async fn send_http_post_request(client: &reqwest::Client, msg: &str, messaging: &MessagingConfig) { if let Some(rest) = &messaging.rest { let data = msg.to_owned(); match client @@ -27,7 +26,7 @@ async fn send_http_post_request(client: &Arc, msg: &str, messag } } -async fn send_telegram_message(client: &Arc, msg: &str, messaging: &MessagingConfig, json: bool) { +async fn send_telegram_message(client: &reqwest::Client, msg: &str, messaging: &MessagingConfig, json: bool) { // TODO use proxy settings if let Some(telegram) = &messaging.telegram { let (message, options) = { @@ -49,7 +48,7 @@ async fn send_telegram_message(client: &Arc, msg: &str, messagi } } -async fn send_pushover_message(client: &Arc, msg: &str, messaging: &MessagingConfig) { +async fn send_pushover_message(client: &reqwest::Client, msg: &str, messaging: &MessagingConfig) { if let Some(pushover) = &messaging.pushover { let encoded_message: String = url::form_urlencoded::Serializer::new(String::new()) .append_pair("token", pushover.token.as_str()) @@ -75,7 +74,7 @@ async fn send_pushover_message(client: &Arc, msg: &str, messagi } } -async fn dispatch_send_message(client: &Arc, kind: MsgKind, cfg: Option<&MessagingConfig>, msg: &str, json: bool) { +async fn dispatch_send_message(client: &reqwest::Client, kind: MsgKind, cfg: Option<&MessagingConfig>, msg: &str, json: bool) { if let Some(messaging) = cfg { if is_enabled(kind, messaging) { tokio::join!( @@ -87,10 +86,10 @@ async fn dispatch_send_message(client: &Arc, kind: MsgKind, cfg } } -pub async fn send_message_json(client: &Arc, kind: MsgKind, cfg: Option<&MessagingConfig>, msg: &str) { +pub async fn send_message_json(client: &reqwest::Client, kind: MsgKind, cfg: Option<&MessagingConfig>, msg: &str) { dispatch_send_message(client, kind, cfg, msg, true).await; } -pub async fn send_message(client: &Arc, kind: MsgKind, cfg: Option<&MessagingConfig>, msg: &str) { +pub async fn send_message(client: &reqwest::Client, kind: MsgKind, cfg: Option<&MessagingConfig>, msg: &str) { dispatch_send_message(client, kind, cfg, msg, false).await; } diff --git a/backend/src/model/xmltv.rs b/backend/src/model/xmltv.rs index d4a5f131c..8954f41f2 100644 --- a/backend/src/model/xmltv.rs +++ b/backend/src/model/xmltv.rs @@ -197,8 +197,9 @@ pub async fn parse_xmltv_for_web_ui_from_file(path: &Path) -> Result, url: &str) -> Result { if let Ok(request_url) = Url::parse(url) { + let client = app_state.http_client.load(); match get_remote_content_as_stream( - Arc::clone(&app_state.http_client.load()), + client.as_ref(), &request_url, InputFetchMethod::GET, None, diff --git a/backend/src/processing/playlist_watch.rs b/backend/src/processing/playlist_watch.rs index d619f7f5c..a67033c6b 100644 --- a/backend/src/processing/playlist_watch.rs +++ b/backend/src/processing/playlist_watch.rs @@ -1,6 +1,5 @@ use std::collections::BTreeSet; use std::path::{Path}; -use std::sync::Arc; use log::{error, info}; use shared::model::{MsgKind, PlaylistGroup}; use crate::messaging::{send_message}; @@ -8,7 +7,7 @@ use crate::model::Config; use crate::utils; use crate::utils::{bincode_deserialize, bincode_serialize}; -pub async fn process_group_watch(client: &Arc, cfg: &Config, target_name: &str, pl: &PlaylistGroup) { +pub async fn process_group_watch(client: &reqwest::Client, cfg: &Config, target_name: &str, pl: &PlaylistGroup) { let mut new_tree = BTreeSet::new(); pl.channels.iter().for_each(|chan| { let header = &chan.header; @@ -60,7 +59,7 @@ struct WatchChanges { pub removed: Vec, } -async fn handle_watch_notification(client: &Arc, cfg: &Config, added: &BTreeSet, removed: &BTreeSet, target_name: &str, group_name: &str) { +async fn handle_watch_notification(client: &reqwest::Client, cfg: &Config, added: &BTreeSet, removed: &BTreeSet, target_name: &str, group_name: &str) { let added = added.iter().map(std::string::ToString::to_string).collect::>(); let removed = removed.iter().map(std::string::ToString::to_string).collect::>(); if !added.is_empty() || !removed.is_empty() { diff --git a/backend/src/processing/processor/playlist.rs b/backend/src/processing/processor/playlist.rs index b168343f8..710bbf190 100644 --- a/backend/src/processing/processor/playlist.rs +++ b/backend/src/processing/processor/playlist.rs @@ -29,7 +29,6 @@ use crate::utils::StepMeasure; use deunicode::deunicode; use futures::StreamExt; use log::{debug, error, info, log_enabled, trace, warn, Level}; -use reqwest::Client; use shared::error::{get_errors_notify_message, notify_err, TuliproxError}; use shared::foundation::filter::{get_field_value, set_field_value, Filter, ValueAccessor, ValueProvider}; use shared::model::{CounterModifier, FieldGetAccessor, FieldSetAccessor, InputType, ItemField, MsgKind, PlaylistEntry, PlaylistGroup, PlaylistItem, PlaylistUpdateState, ProcessingOrder, UUIDType, XtreamCluster}; @@ -305,7 +304,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 playlist_download_from_input(client: &Arc, config: &Arc, input: &Arc) -> (Vec, Vec) { +async fn playlist_download_from_input(client: &reqwest::Client, config: &Arc, input: &Arc) -> (Vec, Vec) { let working_dir = &config.working_dir; match input.input_type { InputType::M3u => m3u::get_m3u_playlist(client, config, input, working_dir).await, @@ -314,7 +313,7 @@ async fn playlist_download_from_input(client: &Arc, config: &Ar } } -async fn process_source(client: Arc, cfg: Arc, source_idx: usize, +async fn process_source(client: &reqwest::Client, cfg: Arc, source_idx: usize, user_targets: Arc, event_manager: Option>, playlist_state: Option<&Arc>, ) -> (Vec, Vec, Vec) { @@ -332,9 +331,9 @@ async fn process_source(client: Arc, cfg: Arc, sourc let working_dir = &config.working_dir; source_downloaded = true; let start_time = Instant::now(); - let (mut playlistgroups, mut error_list) = playlist_download_from_input(&client, &config, input).await; + let (mut playlistgroups, mut error_list) = playlist_download_from_input(client, &config, input).await; let (tvguide, mut tvguide_errors) = if error_list.is_empty() { - epg::get_xmltv(Arc::clone(&client), input, working_dir).await + epg::get_xmltv(client, input, working_dir).await } else { (None, vec![]) }; @@ -373,7 +372,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, playlist_state).await { + match process_playlist_for_target(&cfg, client, &mut source_playlists, target, &mut input_stats, &mut errors, event_manager_clone, playlist_state).await { Ok(()) => { target_stats.push(TargetStats::success(&target.name)); } @@ -407,7 +406,7 @@ fn create_input_stat(group_count: usize, channel_count: usize, error_count: usiz } } -async fn process_sources(client: Arc, config: &Arc, user_targets: Arc, +async fn process_sources(client: &reqwest::Client, config: &Arc, user_targets: Arc, event_manager: Option>, playlist_state: Option<&Arc>, ) -> (Vec, Vec) { let mut async_tasks = JoinSet::new(); @@ -432,13 +431,13 @@ async fn process_sources(client: Arc, config: &Arc, let usr_trgts = user_targets.clone(); let event_manager = event_manager.clone(); if process_parallel { - let http_client = Arc::clone(&client); + let http_client = client.clone(); let playlist_state = playlist_state.cloned(); async_tasks.spawn(async move { // Hold the per-source lock for the full duration of this update. let current_update_lock = update_lock; let (input_stats, target_stats, mut res_errors) = - process_source(Arc::clone(&http_client), cfg, index, usr_trgts, event_manager, playlist_state.as_ref()).await; + process_source(&http_client, cfg, index, usr_trgts, event_manager, playlist_state.as_ref()).await; shared_errors.lock().await.append(&mut res_errors); if let Some(process_stats) = SourceStats::try_new(input_stats, target_stats) { shared_stats.lock().await.push(process_stats); @@ -447,7 +446,7 @@ async fn process_sources(client: Arc, config: &Arc, }); } else { let (input_stats, target_stats, mut res_errors) = - process_source(Arc::clone(&client), cfg, index, usr_trgts, event_manager, playlist_state).await; + process_source(client, cfg, index, usr_trgts, event_manager, playlist_state).await; shared_errors.lock().await.append(&mut res_errors); if let Some(process_stats) = SourceStats::try_new(input_stats, target_stats) { shared_stats.lock().await.push(process_stats); @@ -531,7 +530,7 @@ fn flatten_groups(playlistgroups: Vec) -> Vec { #[allow(clippy::too_many_arguments)] async fn process_playlist_for_target(app_config: &AppConfig, - client: Arc, + client: &reqwest::Client, playlists: &mut [FetchedPlaylist<'_>], target: &ConfigTarget, stats: &mut HashMap, @@ -558,8 +557,8 @@ async fn process_playlist_for_target(app_config: &AppConfig, let mut step = StepMeasure::new(&target.name, broadcast_step); for provider_fpl in playlists.iter_mut() { let mut processed_fpl = execute_pipe(target, &pipe, provider_fpl, &mut duplicates); - playlist_resolve_series(app_config, Arc::clone(&client), target, errors, &pipe, provider_fpl, &mut processed_fpl).await; - playlist_resolve_vod(app_config, Arc::clone(&client), target, errors, &mut processed_fpl).await; + playlist_resolve_series(app_config, client, target, errors, &pipe, provider_fpl, &mut processed_fpl).await; + playlist_resolve_vod(app_config, client, target, errors, &mut processed_fpl).await; // stats let input_stats = stats.get_mut(&processed_fpl.input.name); if let Some(stat) = input_stats { @@ -580,7 +579,7 @@ async fn process_playlist_for_target(app_config: &AppConfig, Ok(()) } else { // Process Trakt categories - if trakt_playlist(&client, target, errors, &mut new_playlist).await { + if trakt_playlist(client, target, errors, &mut new_playlist).await { step.tick("trakt categories"); } @@ -596,7 +595,7 @@ async fn process_playlist_for_target(app_config: &AppConfig, step.tick("assigning channel counter"); let config = app_config.config.load(); - if process_watch(&config, &client, target, &flat_new_playlist).await { + if process_watch(&config, client, target, &flat_new_playlist).await { step.tick("group watches"); } let result = persist_playlist(app_config, &mut flat_new_playlist, flatten_tvguide(&new_epg).as_ref(), target, playlist_state).await; @@ -605,8 +604,8 @@ async fn process_playlist_for_target(app_config: &AppConfig, } } -async fn trakt_playlist(client: &Arc, target: &ConfigTarget, errors: &mut Vec, playlist: &mut Vec) -> bool { - match process_trakt_categories_for_target(Arc::clone(client), playlist, target).await { +async fn trakt_playlist(client: &reqwest::Client, target: &ConfigTarget, errors: &mut Vec, playlist: &mut Vec) -> bool { + match process_trakt_categories_for_target(client, playlist, target).await { Ok(Some(trakt_categories)) => { if !trakt_categories.is_empty() { info!("Adding {} Trakt categories to playlist", trakt_categories.len()); @@ -637,7 +636,7 @@ async fn process_epg(processed_fetched_playlists: &mut Vec>) (new_epg, new_playlist) } -async fn process_watch(cfg: &Config, client: &Arc, target: &ConfigTarget, new_playlist: &[PlaylistGroup]) -> bool { +async fn process_watch(cfg: &Config, client: &reqwest::Client, target: &ConfigTarget, new_playlist: &[PlaylistGroup]) -> bool { if let Some(watches) = &target.watch { if default_as_default().eq_ignore_ascii_case(&target.name) { error!("can't watch a target with no unique name"); @@ -657,10 +656,10 @@ async fn process_watch(cfg: &Config, client: &Arc, target: &Con } } -pub async fn exec_processing(client: Arc, app_config: Arc, targets: Arc, event_manager: Option>, playlist_state: Option>) { +pub async fn exec_processing(client: &reqwest::Client, 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, playlist_state.as_ref()).await; + let (stats, errors) = process_sources(client, &app_config, targets.clone(), event_manager_clone, playlist_state.as_ref()).await; // log errors for err in &errors { error!("{}", err.message); @@ -677,7 +676,7 @@ pub async fn exec_processing(client: Arc, app_config: Arc error!("Failed to serialize playlist stats {err}"), } @@ -692,7 +691,7 @@ pub async fn exec_processing(client: Arc, app_config: Arc Option<&str> { @@ -243,8 +242,8 @@ pub struct TraktCategoriesProcessor { } impl TraktCategoriesProcessor { - pub fn new(http_client: Arc, trakt_config: &TraktConfig) -> Self { - let client = TraktClient::new(http_client, trakt_config.api.clone()); + pub fn new(http_client: &reqwest::Client, trakt_config: &TraktConfig) -> Self { + let client = TraktClient::new(http_client.clone(), trakt_config.api.clone()); Self { client } } @@ -292,7 +291,7 @@ impl TraktCategoriesProcessor { } } pub async fn process_trakt_categories_for_target( - http_client: Arc, + http_client: &reqwest::Client, playlist: &[PlaylistGroup], target: &ConfigTarget, ) -> Result>, Vec> { diff --git a/backend/src/processing/processor/xtream.rs b/backend/src/processing/processor/xtream.rs index f7591d133..9d18c423d 100644 --- a/backend/src/processing/processor/xtream.rs +++ b/backend/src/processing/processor/xtream.rs @@ -9,7 +9,6 @@ use serde::{Deserialize, Serialize}; use std::collections::HashMap; use std::fs::File; use std::path::PathBuf; -use std::sync::Arc; use crate::repository::bplustree::BPlusTree; use crate::repository::storage_const; use crate::repository::xtream_repository::xtream_get_record_file_path; @@ -17,7 +16,7 @@ use crate::utils; use crate::utils::xtream; use serde_json::{from_str, to_string, Value}; -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 { +pub(in crate::processing) async fn playlist_resolve_download_playlist_item(client: &reqwest::Client, 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) = xtream::get_xtream_player_api_info_url(input, cluster, provider_id) { diff --git a/backend/src/processing/processor/xtream_series.rs b/backend/src/processing/processor/xtream_series.rs index 9389cb6e8..6a788132b 100644 --- a/backend/src/processing/processor/xtream_series.rs +++ b/backend/src/processing/processor/xtream_series.rs @@ -14,7 +14,6 @@ use crate::processing::processor::{handle_error, handle_error_and_return, create use std::collections::{HashMap, HashSet}; use std::fs::File; use std::io::{BufWriter, Write}; -use std::sync::Arc; use std::time::Instant; use log::{error, info, log_enabled, warn, Level}; use crate::model::{XtreamSeriesEpisode, XtreamSeriesInfoEpisode}; @@ -50,7 +49,7 @@ fn should_update_series_info(pli: &mut PlaylistItem, processed_provider_ids: &Ha should_update_info(pli, processed_provider_ids, crate::model::XC_TAG_SERIES_INFO_LAST_MODIFIED) } -async fn playlist_resolve_series_info(cfg: &AppConfig, client: Arc, errors: &mut Vec, +async fn playlist_resolve_series_info(cfg: &AppConfig, client: &reqwest::Client, errors: &mut Vec, fpl: &mut FetchedPlaylist<'_>, resolve_delay: u16) -> bool { let mut processed_info_ids = read_processed_series_info_ids(cfg, errors, fpl).await; let mut fetched_in_run: HashSet = HashSet::new(); @@ -83,7 +82,7 @@ async fn playlist_resolve_series_info(cfg: &AppConfig, client: Arc, + client: &reqwest::Client, target: &ConfigTarget, errors: &mut Vec, pipe: &ProcessingPipe, diff --git a/backend/src/processing/processor/xtream_vod.rs b/backend/src/processing/processor/xtream_vod.rs index 865c1b506..8880fb509 100644 --- a/backend/src/processing/processor/xtream_vod.rs +++ b/backend/src/processing/processor/xtream_vod.rs @@ -8,11 +8,9 @@ use crate::repository::xtream_repository::{write_vod_info_to_wal_file, xtream_up use shared::error::{notify_err}; use crate::processing::processor::{handle_error, handle_error_and_return, create_resolve_options_function_for_xtream_target}; use shared::utils::{get_u32_from_serde_value, get_u64_from_serde_value, get_string_from_serde_value}; -use crate::repository::xtream_repository::xtream_get_input_info; use serde_json::{from_str, Map, Value}; use std::collections::{HashMap, HashSet}; use std::io::{Write}; -use std::sync::Arc; use std::time::Instant; use log::{info, log_enabled, Level}; use crate::utils; @@ -65,7 +63,9 @@ fn should_update_vod_info(pli: &mut PlaylistItem, processed_provider_ids: &HashM should_update_info(pli, processed_provider_ids, crate::model::XC_TAG_VOD_INFO_ADDED) } -pub async fn playlist_resolve_vod(app_config: &AppConfig, client: Arc, target: &ConfigTarget, errors: &mut Vec, fpl: &mut FetchedPlaylist<'_>) { +const FLUSH_INTERVAL: usize = 50; + +pub async fn playlist_resolve_vod(app_config: &AppConfig, client: &reqwest::Client, target: &ConfigTarget, errors: &mut Vec, fpl: &mut FetchedPlaylist<'_>) { let (resolve_movies, resolve_delay) = get_resolve_vod_options(target, fpl); if !resolve_movies { return; } @@ -76,7 +76,7 @@ pub async fn playlist_resolve_vod(app_config: &AppConfig, client: Arc = read_processed_vod_info_ids(app_config, errors, fpl).await; let mut fetched_in_run: HashSet = HashSet::new(); let mut content_writer = utils::file_writer(&wal_content_file); let mut record_writer = utils::file_writer(&wal_record_file); @@ -87,78 +87,79 @@ pub async fn playlist_resolve_vod(app_config: &AppConfig, client: Arc>(normalized_str).ok().and_then(|info_doc| { + info_doc.get("info").cloned().map(|info_content| { + let mut wrapped_info = Map::new(); + wrapped_info.insert("info".to_string(), info_content); + Value::Object(wrapped_info) + }) + }); + } + } + } + if log_enabled!(Level::Info) { + processed_vod_info_count += 1; + if last_log_time.elapsed().as_secs() >= 30 { + info!("resolved {processed_vod_info_count}/{vod_info_count} vod info"); + last_log_time = Instant::now(); } } } - if log_enabled!(Level::Info) { - processed_vod_info_count += 1; - let elapsed = start_time.elapsed().as_secs(); - if elapsed > 0 && ((processed_vod_info_count - last_processed_vod_info_count) > 50) && elapsed.is_multiple_of(30) { - info!("resolved {processed_vod_info_count}/{vod_info_count} vod info"); - last_processed_vod_info_count = processed_vod_info_count; - } - } - } - if last_processed_vod_info_count != processed_vod_info_count { - info!("resolved {processed_vod_info_count}/{vod_info_count} vod info"); } + info!("resolved {processed_vod_info_count}/{vod_info_count} vod info"); if content_updated { // TODO better approach for transactional updates is multiplexed WAL file. - + // final flush & sync with proper error handling handle_error!(content_writer.flush(), |err| errors.push(notify_err!(format!("Failed to resolve vod, could not write to wal file {err}")))); handle_error!(record_writer.flush(), |err| errors.push(notify_err!(format!("Failed to resolve vod tmdb, could not write to wal file {err}")))); handle_error!(content_writer.get_ref().sync_all(), |err| errors.push(notify_err!(format!("Failed to sync vod info to wal file {err}")))); handle_error!(record_writer.get_ref().sync_all(), |err| errors.push(notify_err!(format!("Failed to sync vod info record to wal file {err}")))); + // drop writers and files to release handles drop(content_writer); drop(record_writer); drop(wal_content_file); drop(wal_record_file); + handle_error!(xtream_update_input_info_file(app_config, fpl.input, &wal_content_path, XtreamCluster::Video).await, |err| errors.push(err)); handle_error!(xtream_update_input_vod_record_from_wal_file(app_config, fpl.input, &wal_record_path).await, |err| errors.push(err)); } - - // Update in-memory playlist items with the newly fetched vod info. - // This makes the data available for subsequent processing steps like STRM export. - let vod_info_iter = fpl.playlistgroups.iter_mut() - .flat_map(|plg| &mut plg.channels) - .filter(|pli| pli.header.xtream_cluster == XtreamCluster::Video); - - for pli in vod_info_iter { - if let Some(provider_id) = pli.header.get_provider_id() { - if let Some(content) = xtream_get_input_info(app_config, fpl.input, provider_id, XtreamCluster::Video).await { - // Add the "info" section to the playlist item additional properties. - pli.header.additional_properties = from_str::>(&content).ok().and_then(|info_doc| { - info_doc.get("info").cloned().map(|info_content| { - let mut wrapped_info = Map::new(); - wrapped_info.insert("info".to_string(), info_content); - Value::Object(wrapped_info) - }) - }); - } - } - } } diff --git a/backend/src/utils/network/epg.rs b/backend/src/utils/network/epg.rs index 2dbe95ee9..1e51b8a90 100644 --- a/backend/src/utils/network/epg.rs +++ b/backend/src/utils/network/epg.rs @@ -5,7 +5,6 @@ use crate::utils::{add_prefix_to_filename, cleanup_unlisted_files_with_suffix, p use crate::utils::request; use log::debug; use std::path::PathBuf; -use std::sync::Arc; use shared::utils::{sanitize_sensitive_info, short_hash}; fn get_input_raw_epg_file_path(url: &str, input: &ConfigInput, working_dir: &str) -> Option { @@ -14,13 +13,13 @@ fn get_input_raw_epg_file_path(url: &str, input: &ConfigInput, working_dir: &str .map(|path| add_prefix_to_filename(&path, format!("{file_prefix}_epg_").as_str(), Some("xml"))) } -async fn download_epg_file(url: &str, client: &Arc, input: &ConfigInput, working_dir: &str) -> Result { +async fn download_epg_file(url: &str, client: &reqwest::Client, input: &ConfigInput, working_dir: &str) -> Result { debug!("Getting epg file path for url: {}", sanitize_sensitive_info(url)); let persist_file_path = get_input_raw_epg_file_path(url, input, working_dir); - request::get_input_epg_content_as_file(Arc::clone(client), input, working_dir, url, persist_file_path).await + request::get_input_epg_content_as_file(client, input, working_dir, url, persist_file_path).await } -pub async fn get_xmltv(client: Arc, input: &ConfigInput, working_dir: &str) -> (Option, Vec) { +pub async fn get_xmltv(client: &reqwest::Client, input: &ConfigInput, working_dir: &str) -> (Option, Vec) { match &input.epg { None => (None, vec![]), Some(epg_config) => { @@ -29,7 +28,7 @@ pub async fn get_xmltv(client: Arc, input: &ConfigInput, workin let mut stored_file_paths = vec![]; for epg_source in &epg_config.sources { - match download_epg_file(&epg_source.url, &client, input, working_dir).await { + match download_epg_file(&epg_source.url, client, input, working_dir).await { Ok(file_path) => { stored_file_paths.push(file_path.clone()); file_paths.push(PersistedEpgSource {file_path, priority: epg_source.priority, logo_override: epg_source.logo_override}); diff --git a/backend/src/utils/network/ip_checker.rs b/backend/src/utils/network/ip_checker.rs index 230f5a235..1ee5dcdf4 100644 --- a/backend/src/utils/network/ip_checker.rs +++ b/backend/src/utils/network/ip_checker.rs @@ -1,11 +1,9 @@ use shared::error::{TuliproxError, TuliproxErrorKind}; use crate::model::IpCheckConfig; use regex::Regex; -use reqwest::Client; -use std::sync::Arc; use shared::utils::sanitize_sensitive_info; -async fn fetch_ip(client: &Arc, url: &str, regex: Option<&Regex>) -> Result { +async fn fetch_ip(client: &reqwest::Client, url: &str, regex: Option<&Regex>) -> Result { let response = client.get(url).send().await .map_err(|e| TuliproxError::new(TuliproxErrorKind::Info, format!("Failed to request {}: {e}", sanitize_sensitive_info(url))))?; @@ -28,7 +26,7 @@ async fn fetch_ip(client: &Arc, url: &str, regex: Option<&Regex>) -> Res } /// Fetch both IPs from a shared URL (if both regex patterns are available) -async fn fetch_combined_ips(client: &Arc, config: &IpCheckConfig, url: &str) -> (Option, Option) { +async fn fetch_combined_ips(client: &reqwest::Client, config: &IpCheckConfig, url: &str) -> (Option, Option) { let response = client.get(url).send().await.ok(); let text = match response { Some(r) => r.text().await.ok(), @@ -55,7 +53,7 @@ async fn fetch_combined_ips(client: &Arc, config: &IpCheckConfig, url: & } /// Fetch both IPv4 and IPv6 addresses, using separate or combined URL(s) -pub async fn get_ips(client: &Arc, config: &IpCheckConfig) -> Result<(Option, Option), TuliproxError> { +pub async fn get_ips(client: &reqwest::Client, config: &IpCheckConfig) -> Result<(Option, Option), TuliproxError> { match (&config.url_ipv4, &config.url_ipv6, &config.url) { // Both dedicated URLs provided (Some(url_v4), Some(url_v6), _) => { diff --git a/backend/src/utils/network/m3u.rs b/backend/src/utils/network/m3u.rs index a29e0351a..d640d177e 100644 --- a/backend/src/utils/network/m3u.rs +++ b/backend/src/utils/network/m3u.rs @@ -6,7 +6,7 @@ use crate::processing::parser::m3u; use crate::utils::prepare_file_path; use crate::utils::request; -pub async fn get_m3u_playlist(client: &Arc, cfg: &Arc, input: &Arc, working_dir: &str) -> (Vec, Vec) { +pub async fn get_m3u_playlist(client: &reqwest::Client, cfg: &Arc, input: &Arc, working_dir: &str) -> (Vec, Vec) { let input_source: InputSource = { match input.staged.as_ref() { None => input.as_ref().into(), diff --git a/backend/src/utils/network/request.rs b/backend/src/utils/network/request.rs index 8c907fd28..4b15cde7b 100644 --- a/backend/src/utils/network/request.rs +++ b/backend/src/utils/network/request.rs @@ -2,7 +2,6 @@ use std::collections::{HashMap, HashSet}; use std::io::{Error, ErrorKind}; use std::path::{Path, PathBuf}; use std::pin::Pin; -use std::sync::Arc; use std::time::{Duration, Instant}; use futures::{StreamExt, TryStreamExt}; use log::{debug, error, log_enabled, trace, Level}; @@ -53,7 +52,7 @@ pub fn classify_content_type(headers: &[(String, String)]) -> MimeCategory { }) } -pub async fn get_input_epg_content_as_file(client: Arc, input: &ConfigInput, working_dir: &str, url_str: &str, persist_filepath: Option) -> Result { +pub async fn get_input_epg_content_as_file(client: &reqwest::Client, input: &ConfigInput, working_dir: &str, url_str: &str, persist_filepath: Option) -> Result { debug_if_enabled!("getting input epg content working_dir: {}, url: {}", working_dir, sanitize_sensitive_info(url_str)); if url_str.parse::().is_ok() { match download_epg_content_as_file(client, input, url_str, working_dir, persist_filepath).await { @@ -96,11 +95,11 @@ pub async fn get_input_epg_content_as_file(client: Arc, input: } -pub async fn get_input_text_content(client: &Arc, input: &InputSource, working_dir: &str, persist_filepath: Option) -> Result { +pub async fn get_input_text_content(client: &reqwest::Client, input: &InputSource, working_dir: &str, persist_filepath: Option) -> Result { debug_if_enabled!("getting input text content working_dir: {}, url: {}", working_dir, sanitize_sensitive_info(&input.url)); if input.url.parse::().is_ok() { - match download_text_content(Arc::clone(client), None, input, None, persist_filepath).await { + match download_text_content(client, None, input, None, persist_filepath).await { Ok((content, _response_url)) => Ok(content), Err(e) => { error!("Failed to download input '{}': {}", &input.name, sanitize_sensitive_info(e.to_string().as_str())); @@ -140,7 +139,7 @@ pub async fn get_input_text_content(client: &Arc, input: &Input } pub fn get_client_request - (client: &Arc, + (client: &reqwest::Client, method: InputFetchMethod, headers: Option<&HashMap>, url: &Url, @@ -242,9 +241,9 @@ pub async fn get_local_file_content(file_path: &Path) -> Result, input: &ConfigInput, url: &Url, file_path: &Path) -> Result { +async fn get_remote_content_as_file(client: &reqwest::Client, input: &ConfigInput, url: &Url, file_path: &Path) -> Result { let start_time = Instant::now(); - let request = get_client_request(&client, input.method, Some(&input.headers), url, None, None); + let request = get_client_request(client, input.method, Some(&input.headers), url, None, None); match request.send().await { Ok(response) => { if response.status().is_success() { @@ -279,12 +278,12 @@ type DynReader = Pin>; #[allow(clippy::implicit_hasher)] pub async fn get_remote_content_as_stream( - client: Arc, + client: &reqwest::Client, url: &Url, method: InputFetchMethod, headers: Option<&HashMap> ) -> Result<(DynReader, String), Error> { - let request = get_client_request(&client, method, headers, url, None, None); + let request = get_client_request(client, method, headers, url, None, None); let response = request.send().await.map_err(std::io::Error::other)?; if !response.status().is_success() { @@ -320,7 +319,7 @@ pub async fn get_remote_content_as_stream( Ok((reader, response_url)) } -async fn get_remote_content(client: Arc, input: &InputSource, headers: Option<&HeaderMap>, url: &Url, disabled_headers: Option<&ReverseProxyDisabledHeaderConfig>) -> Result<(String, String), Error> { +async fn get_remote_content(client: &reqwest::Client, input: &InputSource, headers: Option<&HeaderMap>, url: &Url, disabled_headers: Option<&ReverseProxyDisabledHeaderConfig>) -> Result<(String, String), Error> { let start_time = Instant::now(); let custom_headers = headers.map(|h| { @@ -328,14 +327,14 @@ async fn get_remote_content(client: Arc, input: &InputSource, h let merged = get_request_headers(Some(&input.headers), custom_headers.as_ref(), disabled_headers); let headers: HashMap = merged.iter().map(|(k, v)| (k.as_str().to_string(), String::from_utf8_lossy(v.as_bytes()).to_string())).collect(); - let (mut stream, response_url) = get_remote_content_as_stream(client.clone(), url, input.method, Some(&headers)).await.map_err(|e| str_to_io_error(&format!("Failed to read content: {e}")))?; + let (mut stream, response_url) = get_remote_content_as_stream(client, url, input.method, Some(&headers)).await.map_err(|e| str_to_io_error(&format!("Failed to read content: {e}")))?; let mut content = String::new(); stream.read_to_string(&mut content).await.map_err(|e| str_to_io_error(&format!("Failed to read content: {e}")))?; debug_if_enabled!("Request took: {} {}", format_elapsed_time(start_time.elapsed().as_secs()), sanitize_sensitive_info(url.as_str())); Ok((content, response_url)) } -async fn download_epg_content_as_file(client: Arc, input: &ConfigInput, url_str: &str, working_dir: &str, persist_filepath: Option) -> Result { +async fn download_epg_content_as_file(client: &reqwest::Client, 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 {}", sanitize_sensitive_info(url_str)))), |file_path| if file_path.exists() { @@ -361,7 +360,7 @@ async fn download_epg_content_as_file(client: Arc, input: &Conf } pub async fn download_text_content( - client: Arc, + client: &reqwest::Client, disabled_headers: Option<&ReverseProxyDisabledHeaderConfig>, input: &InputSource, headers: Option<&HeaderMap>, @@ -396,7 +395,7 @@ pub async fn download_text_content( } } -async fn download_json_content(client: Arc, disabled_headers: Option<&ReverseProxyDisabledHeaderConfig>, input: &InputSource, persist_filepath: Option) -> Result { +async fn download_json_content(client: &reqwest::Client, disabled_headers: Option<&ReverseProxyDisabledHeaderConfig>, input: &InputSource, persist_filepath: Option) -> Result { debug_if_enabled!("downloading json content from {}", sanitize_sensitive_info(&input.url)); match download_text_content(client, disabled_headers, input, None, persist_filepath).await { Ok((content, _response_url)) => { @@ -409,7 +408,7 @@ async fn download_json_content(client: Arc, disabled_headers: O } } -pub async fn get_input_json_content(client: Arc, disabled_headers: Option<&ReverseProxyDisabledHeaderConfig>, input: &InputSource, persist_filepath: Option) -> Result { +pub async fn get_input_json_content(client: &reqwest::Client, disabled_headers: Option<&ReverseProxyDisabledHeaderConfig>, input: &InputSource, persist_filepath: Option) -> Result { match download_json_content(client, disabled_headers, input, persist_filepath).await { Ok(content) => Ok(content), Err(e) => create_tuliprox_error_result!(TuliproxErrorKind::Notify, "cant download input {}, url: {} => {}", input.name, sanitize_sensitive_info(&input.url), sanitize_sensitive_info(e.to_string().as_str())) diff --git a/backend/src/utils/network/xtream.rs b/backend/src/utils/network/xtream.rs index 2e126b40e..005345bd0 100644 --- a/backend/src/utils/network/xtream.rs +++ b/backend/src/utils/network/xtream.rs @@ -46,7 +46,7 @@ pub fn get_xtream_player_api_info_url(input: &ConfigInput, cluster: XtreamCluste } -pub async fn get_xtream_stream_info_content(client: Arc, input: &InputSource) -> Result { +pub async fn get_xtream_stream_info_content(client: &reqwest::Client, input: &InputSource) -> Result { match request::download_text_content(client, None, input, None, None).await { Ok((content, _response_url)) => Ok(content), Err(err) => Err(err) @@ -54,7 +54,7 @@ pub async fn get_xtream_stream_info_content(client: Arc, input: } #[allow(clippy::too_many_arguments)] -pub async fn get_xtream_stream_info

(client: Arc, +pub async fn get_xtream_stream_info

(client: &reqwest::Client, app_state: &Arc, user: &ProxyUserCredentials, input: &ConfigInput, @@ -135,12 +135,12 @@ const ACTIONS: [(XtreamCluster, &str, &str); 3] = [ (XtreamCluster::Video, crate::model::XC_ACTION_GET_VOD_CATEGORIES, crate::model::XC_ACTION_GET_VOD_STREAMS), (XtreamCluster::Series, crate::model::XC_ACTION_GET_SERIES_CATEGORIES, crate::model::XC_ACTION_GET_SERIES)]; -async fn xtream_login(cfg: &Config, client: &Arc, input: &InputSource, username: &str) -> Result, TuliproxError> { - let content = if let Ok(content) = request::get_input_json_content(Arc::clone(client), None, input, None).await { +async fn xtream_login(cfg: &Config, client: &reqwest::Client, input: &InputSource, username: &str) -> Result, TuliproxError> { + let content = if let Ok(content) = request::get_input_json_content(client, None, input, None).await { content } else { let input_source_account_info = input.with_url(format!("{}&action={}", &input.url, crate::model::XC_ACTION_GET_ACCOUNT_INFO)); - match request::get_input_json_content(Arc::clone(client), None, &input_source_account_info, None).await { + match request::get_input_json_content(client, None, &input_source_account_info, None).await { Ok(content) => content, Err(err) => { warn!("Failed to login xtream account {username} {err}"); @@ -182,7 +182,7 @@ async fn xtream_login(cfg: &Config, client: &Arc, input: &Input } } -pub async fn notify_account_expire(exp_date: Option, cfg: &Config, client: &Arc, username: &str, input_name: &str) { +pub async fn notify_account_expire(exp_date: Option, cfg: &Config, client: &reqwest::Client, username: &str, input_name: &str) { if let Some(expiration_timestamp) = exp_date { let now_secs = Utc::now().timestamp(); // UTC-Time if expiration_timestamp > now_secs { @@ -202,7 +202,7 @@ pub async fn notify_account_expire(exp_date: Option, cfg: &Config, client: } } -pub async fn get_xtream_playlist(cfg: &Arc, client: &Arc, input: &Arc, working_dir: &str) -> (Vec, Vec) { +pub async fn get_xtream_playlist(cfg: &Arc, client: &reqwest::Client, input: &Arc, working_dir: &str) -> (Vec, Vec) { let input_source: InputSource = { match input.staged.as_ref() { None => input.as_ref().into(), @@ -235,8 +235,8 @@ pub async fn get_xtream_playlist(cfg: &Arc, client: &Arc { match xtream::parse_xtream(input, @@ -265,7 +265,7 @@ pub async fn get_xtream_playlist(cfg: &Arc, client: &Arc, client: &Arc, input: &Arc) { +async fn check_alias_user_state(cfg: &Arc, client: &reqwest::Client, input: &Arc) { if let Some(aliases) = input.aliases.as_ref() { for alias in aliases { if is_input_expired(alias.exp_date) { diff --git a/backend/src/utils/telegram.rs b/backend/src/utils/telegram.rs index aaae5c008..4755775fb 100644 --- a/backend/src/utils/telegram.rs +++ b/backend/src/utils/telegram.rs @@ -1,4 +1,3 @@ -use std::sync::Arc; use log::{debug, error}; use url::Url; @@ -63,7 +62,7 @@ pub fn telegram_create_instance(bot_token: &str, chat_id: &str) -> BotInstance { } pub async fn telegram_send_message( - client: &Arc, + client: &reqwest::Client, instance: &BotInstance, msg: &str, options: Option<&SendMessageOption>, diff --git a/backend/src/utils/trakt/client.rs b/backend/src/utils/trakt/client.rs index 43adf51e6..ae2c473ee 100644 --- a/backend/src/utils/trakt/client.rs +++ b/backend/src/utils/trakt/client.rs @@ -1,21 +1,20 @@ use crate::model::{TraktApiConfig, TraktListConfig, TraktListItem}; use shared::error::TuliproxError; use reqwest::header::{HeaderMap, HeaderValue}; -use std::sync::Arc; use log::{debug, info}; use shared::error::{info_err}; use shared::utils::trim_last_slash; use super::errors::{handle_trakt_api_error}; pub struct TraktClient { - client: Arc, + client: reqwest::Client, api_config: TraktApiConfig, // Pre-computed headers to avoid recreating them each time headers: HeaderMap, } impl TraktClient { - pub fn new(client: Arc, api_config: TraktApiConfig) -> Self { + pub fn new(client: reqwest::Client, api_config: TraktApiConfig) -> Self { let headers = Self::create_headers(&api_config); Self { client,