From ccdbf1dfce946ecf7f8bb07db67f0f92be43e492 Mon Sep 17 00:00:00 2001 From: euzu Date: Thu, 22 Jan 2026 13:30:59 +0100 Subject: [PATCH] Retry request from reverseproxy retry options --- backend/src/api/api_utils.rs | 143 ++---- .../src/api/endpoints/api_playlist_utils.rs | 2 +- backend/src/api/endpoints/download_api.rs | 5 +- backend/src/api/endpoints/hls_api.rs | 3 +- backend/src/api/endpoints/v1_api.rs | 5 +- backend/src/api/endpoints/v1_api_config.rs | 5 +- backend/src/api/endpoints/v1_api_playlist.rs | 4 +- backend/src/api/endpoints/xtream_api.rs | 21 +- backend/src/api/main_api.rs | 3 +- backend/src/api/model/app_state.rs | 8 +- backend/src/api/scheduler.rs | 3 +- backend/src/main.rs | 2 +- backend/src/model/config/app.rs | 8 +- backend/src/model/config/base.rs | 18 +- backend/src/model/xmltv.rs | 6 +- backend/src/processing/processor/playlist.rs | 11 +- backend/src/processing/processor/xtream.rs | 7 +- .../src/processing/processor/xtream_series.rs | 3 +- .../src/processing/processor/xtream_vod.rs | 3 +- backend/src/utils/network/epg.rs | 31 +- backend/src/utils/network/m3u.rs | 5 +- backend/src/utils/network/request.rs | 445 +++++++++++------- backend/src/utils/network/xtream.rs | 22 +- shared/src/error/tuliprox_error.rs | 9 +- shared/src/model/short_epg.rs | 11 +- 25 files changed, 414 insertions(+), 369 deletions(-) diff --git a/backend/src/api/api_utils.rs b/backend/src/api/api_utils.rs index bc7d7fb7e..dcaa43a80 100644 --- a/backend/src/api/api_utils.rs +++ b/backend/src/api/api_utils.rs @@ -7,10 +7,10 @@ use crate::api::model::{create_channel_unavailable_stream, create_custom_video_s use crate::api::model::{tee_stream, UserSession}; use crate::api::model::{ProviderAllocation, ProviderConfig, ProviderStreamState, StreamDetails, StreamingStrategy}; use crate::auth::Fingerprint; -use crate::model::{ConfigInput, ResourceRetryConfig}; +use crate::model::{ConfigInput}; use crate::model::{ConfigTarget, ProxyUserCredentials}; use crate::tools::lru_cache::LRUResourceCache; -use crate::utils::request::{content_type_from_ext, parse_range}; +use crate::utils::request::{content_type_from_ext, parse_range, send_with_retry}; use crate::utils::{async_file_reader, async_file_writer, create_new_file_for_write, get_file_extension}; use crate::utils::{debug_if_enabled, trace_if_enabled}; use crate::utils::request; @@ -25,7 +25,6 @@ use chrono::{DateTime, Utc}; use futures::{stream, StreamExt, TryStreamExt}; use jsonwebtoken::{decode, Algorithm, DecodingKey, Validation}; use log::{debug, error, info, log_enabled, trace, warn}; -use reqwest::header::RETRY_AFTER; use serde::Serialize; use shared::concat_string; use shared::model::{Claims, InputFetchMethod, PlaylistEntry, PlaylistItemType, ProxyType, StreamChannel, TargetType, UserConnectionPermission, VirtualId, XtreamCluster}; @@ -39,7 +38,6 @@ use std::convert::Infallible; use std::io::SeekFrom; use std::path::{Path, PathBuf}; use std::sync::Arc; -use std::time::Duration; use tokio::io::{AsyncReadExt, AsyncSeekExt}; use tokio::sync::Mutex; use tokio_util::io::ReaderStream; @@ -1343,102 +1341,59 @@ async fn fetch_resource_with_retry( ) -> Option { let config = app_state.app_config.config.load(); let default_user_agent = config.default_user_agent.clone(); - let (max_attempts, backoff_ms, backoff_multiplier) = config - .reverse_proxy - .as_ref() - .map_or_else(ResourceRetryConfig::get_default_retry_values, |rp| rp.resource_retry.get_retry_values()); drop(config); + let disabled_headers = app_state.get_disabled_headers(); - for attempt in 0..max_attempts { - let client = request::get_client_request( - &app_state.http_client.load(), - input.map_or(InputFetchMethod::GET, |i| i.method), - input.map(|i| &i.headers), - url, - Some(req_headers), - disabled_headers.as_ref(), - default_user_agent.as_deref(), + + let response = match send_with_retry( + &app_state.app_config, + url, + || { + request::get_client_request( + &app_state.http_client.load(), + input.map_or(InputFetchMethod::GET, |i| i.method), + input.map(|i| &i.headers), + url, + Some(req_headers), + disabled_headers.as_ref(), + default_user_agent.as_deref(), + ) + }, + ) + .await + { + Ok(response) => response, + Err(_) => return None, + }; + + let status = response.status(); + + if status.is_success() { + return Some( + build_resource_stream_response(app_state, resource_url, response).await, ); - match client.send().await { - Ok(response) => { - let status = response.status(); - if status.is_success() { - return Some( - build_resource_stream_response(app_state, resource_url, response).await, - ); - } - // Retry only for 408, 425, 429 and all 5xx statuses - let should_retry = status.is_server_error() - || matches!( - status, - // reqwest::StatusCode::BAD_REQUEST // 400 is typically client error; retrying likely won't help and adds load. - reqwest::StatusCode::REQUEST_TIMEOUT - | reqwest::StatusCode::TOO_EARLY - | reqwest::StatusCode::TOO_MANY_REQUESTS - ); - - if attempt < max_attempts - 1 && should_retry { - let wait_dur = response - .headers() - .get(RETRY_AFTER) - .and_then(|h| h.to_str().ok()) - .and_then(|s| s.parse::().ok()) - .map_or_else( - || { - let delay = calculate_retry_backoff(backoff_ms, backoff_multiplier, attempt); - Duration::from_millis(delay) - }, - Duration::from_secs, - ); - tokio::time::sleep(wait_dur).await; - continue; - } - - // For non-retriable statuses or when attempts are exhausted, return upstream response including body - debug_if_enabled!( - "Failed to open resource got status {status} for {}", - sanitize_sensitive_info(resource_url) - ); - let mut response_builder = axum::response::Response::builder().status(status); - for (key, value) in response.headers() { - response_builder = response_builder.header(key, value); - } - let stream = response - .bytes_stream() - .map_err(|err| StreamError::reqwest(&err)); - return Some(try_unwrap_body!( - response_builder.body(axum::body::Body::from_stream(stream)) - )); - } - Err(err) => { - if attempt < max_attempts - 1 { - let delay = calculate_retry_backoff(backoff_ms, backoff_multiplier, attempt); - tokio::time::sleep(Duration::from_millis(delay)).await; - continue; - } - error!("Received failure from server {}: {err}", sanitize_sensitive_info(resource_url)); - } - } - break; } - None + + // Non-retriable Status → Upstream Response incl. Body + debug_if_enabled!( + "Failed to open resource got status {status} for {}", + sanitize_sensitive_info(resource_url) + ); + + let mut response_builder = axum::response::Response::builder().status(status); + for (key, value) in response.headers() { + response_builder = response_builder.header(key, value); + } + + let stream = response + .bytes_stream() + .map_err(|err| StreamError::reqwest(&err)); + + Some(try_unwrap_body!( + response_builder.body(axum::body::Body::from_stream(stream)) + )) } -#[allow(clippy::cast_possible_truncation, clippy::cast_sign_loss, clippy::cast_precision_loss)] -fn calculate_retry_backoff(base_delay_ms: u64, multiplier: f64, attempt: u32) -> u64 { - let base = base_delay_ms.max(1); - if multiplier <= 1.0 { - return base; - } - let delay = (base as f64) * multiplier.powi(i32::try_from(attempt).unwrap_or(i32::MAX)); - if !delay.is_finite() || delay < 1.0 { - base - } else if delay >= u64::MAX as f64 { - u64::MAX - } else { - delay as u64 - } -} /// # Panics pub async fn resource_response( diff --git a/backend/src/api/endpoints/api_playlist_utils.rs b/backend/src/api/endpoints/api_playlist_utils.rs index 037cbc9fe..9c8d14819 100644 --- a/backend/src/api/endpoints/api_playlist_utils.rs +++ b/backend/src/api/endpoints/api_playlist_utils.rs @@ -65,7 +65,7 @@ pub(in crate::api::endpoints) async fn get_playlist_for_custom_provider(client: Some(input) => { let (result, errors) = match input.input_type { - InputType::M3u | InputType::M3uBatch => m3u::download_m3u_playlist(client, &cfg, input).await, + InputType::M3u | InputType::M3uBatch => m3u::download_m3u_playlist(app_config, client, &cfg, input).await, InputType::Xtream | InputType::XtreamBatch => { let (pl, err, _) = xtream::download_xtream_playlist(app_config, client, input, Some(&[cluster])).await; (pl, err) diff --git a/backend/src/api/endpoints/download_api.rs b/backend/src/api/endpoints/download_api.rs index 638efa4d2..c33ca69c8 100644 --- a/backend/src/api/endpoints/download_api.rs +++ b/backend/src/api/endpoints/download_api.rs @@ -85,10 +85,7 @@ async fn run_download_queue(cfg: &AppConfig, download_cfg: &VideoDownloadConfig, if next_download.is_some() { { *download_queue.as_ref().active.write().await = next_download; } let config = cfg.config.load(); - let disabled_headers = config - .reverse_proxy - .as_ref() - .and_then(|r| r.disabled_header.clone()); + let disabled_headers = cfg.get_disabled_headers(); let headers = request::get_request_headers( Some(&download_cfg.headers), None, diff --git a/backend/src/api/endpoints/hls_api.rs b/backend/src/api/endpoints/hls_api.rs index 9ced11844..ae21af8ed 100644 --- a/backend/src/api/endpoints/hls_api.rs +++ b/backend/src/api/endpoints/hls_api.rs @@ -133,13 +133,12 @@ pub(in crate::api) async fn handle_hls_stream_request( ); let input_source = InputSource::from(input).with_url(request_url); match request::download_text_content( + &app_state.app_config, &app_state.http_client.load(), - disabled_headers.as_ref(), &input_source, Some(&headers), None, false, - default_user_agent.as_deref(), ) .await { diff --git a/backend/src/api/endpoints/v1_api.rs b/backend/src/api/endpoints/v1_api.rs index df6c11d16..6d28b9240 100644 --- a/backend/src/api/endpoints/v1_api.rs +++ b/backend/src/api/endpoints/v1_api.rs @@ -90,16 +90,13 @@ async fn geoip_update(axum::extract::State(app_state): axum::extract::State 0 { let input_name = &pli.input_name; @@ -934,15 +934,13 @@ async fn xtream_get_short_epg( // TODO serve epg from own db let input_source = InputSource::from(&*input).with_url(info_url); - let default_user_agent = config.default_user_agent.clone(); return match request::download_text_content( + &app_state.app_config, &app_state.http_client.load(), - None, &input_source, None, None, false, - default_user_agent.as_deref(), ) .await { @@ -1041,13 +1039,12 @@ async fn xtream_get_catchup_response( pli.provider_id ))); let input_source = InputSource::from(&*input).with_url(info_url); - let default_user_agent = app_state.app_config.config.load().default_user_agent.clone(); let content = try_result_bad_request!( xtream::get_xtream_stream_info_content( + &app_state.app_config, &app_state.http_client.load(), &input_source, false, - default_user_agent.as_deref(), ) .await ); diff --git a/backend/src/api/main_api.rs b/backend/src/api/main_api.rs index 27c05222d..3366847cc 100644 --- a/backend/src/api/main_api.rs +++ b/backend/src/api/main_api.rs @@ -132,8 +132,9 @@ fn exec_update_on_boot( let playlist_state = Arc::clone(&app_state.playlists); let client = client.clone(); let update_guard = Some(app_state.update_guard.clone()); + let disabled_headers = app_state.get_disabled_headers(); tokio::spawn(async move { - playlist::exec_processing(&client, app_config_clone, targets_clone, None, Some(playlist_state), update_guard).await; + playlist::exec_processing(&client, app_config_clone, targets_clone, None, Some(playlist_state), update_guard, disabled_headers).await; }); } } diff --git a/backend/src/api/model/app_state.rs b/backend/src/api/model/app_state.rs index 53caa6ac7..127800346 100644 --- a/backend/src/api/model/app_state.rs +++ b/backend/src/api/model/app_state.rs @@ -443,13 +443,7 @@ impl AppState { } pub fn get_disabled_headers(&self) -> Option { - self - .app_config - .config - .load() - .reverse_proxy - .as_ref() - .and_then(|r| r.disabled_header.clone()) + self.app_config.get_disabled_headers() } } diff --git a/backend/src/api/scheduler.rs b/backend/src/api/scheduler.rs index e29677031..95c8cf635 100644 --- a/backend/src/api/scheduler.rs +++ b/backend/src/api/scheduler.rs @@ -62,9 +62,10 @@ async fn start_scheduler(client: reqwest::Client, expression: &str, app_state: A 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(); + let disabled_headers = app_state.get_disabled_headers(); sync_panel_api_exp_dates_on_boot(&app_state).await; exec_processing(&client, app_config, Arc::clone(&targets), Some(event_manager), - Some(playlist_state), Some(app_state.update_guard.clone())).await; + Some(playlist_state), Some(app_state.update_guard.clone()), disabled_headers).await; } () = cancel.cancelled() => { break; diff --git a/backend/src/main.rs b/backend/src/main.rs index c54f23f56..b8cef8ab2 100644 --- a/backend/src/main.rs +++ b/backend/src/main.rs @@ -219,7 +219,7 @@ async fn start_in_cli_mode(cfg: Arc, targets: Arc) { error!("Failed to build client {err}"); reqwest::Client::new() }); - playlist::exec_processing(&client, cfg, targets, None, None, None).await; + playlist::exec_processing(&client, cfg, targets, None, None, None, None).await; } async fn start_in_server_mode(cfg: Arc, targets: Arc) { diff --git a/backend/src/model/config/app.rs b/backend/src/model/config/app.rs index 51adf557a..6b017b479 100644 --- a/backend/src/model/config/app.rs +++ b/backend/src/model/config/app.rs @@ -1,5 +1,6 @@ use crate::api::model::TransportStreamBuffer; -use crate::model::{ApiProxyConfig, ApiProxyServerInfo, Config, ConfigInput, ConfigInputOptions, ConfigTarget, CustomStreamResponse, HdHomeRunConfig, Mappings, ProxyUserCredentials, SourcesConfig, TargetOutput}; +use crate::model::{ApiProxyConfig, ApiProxyServerInfo, Config, ConfigInput, ConfigInputOptions, + ConfigTarget, CustomStreamResponse, HdHomeRunConfig, Mappings, ProxyUserCredentials, ReverseProxyDisabledHeaderConfig, SourcesConfig, TargetOutput}; use crate::utils; use arc_swap::access::Access; use arc_swap::{ArcSwap, ArcSwapOption}; @@ -414,5 +415,10 @@ impl AppConfig { let server_info_name = user.server.as_ref().map_or("default", |server_name| server_name.as_str()); self.get_server_info(server_info_name) } + + pub fn get_disabled_headers(&self) -> Option { + let config = > as Access>::load(&self.config); + config.get_disabled_headers() + } } diff --git a/backend/src/model/config/base.rs b/backend/src/model/config/base.rs index 366101274..99bfab6d9 100644 --- a/backend/src/model/config/base.rs +++ b/backend/src/model/config/base.rs @@ -1,13 +1,13 @@ -use std::borrow::Cow; -use std::path::{Path, PathBuf}; +use crate::model::{macros, ConfigApi, LibraryConfig, ReverseProxyConfig, ReverseProxyDisabledHeaderConfig, ScheduleConfig}; +use crate::model::{HdHomeRunConfig, IpCheckConfig, LogConfig, MessagingConfig, ProxyConfig, VideoConfig, WebUiConfig}; +use crate::utils; use log::{error, info}; use path_clean::PathClean; -use shared::error::{TuliproxError}; +use shared::error::TuliproxError; use shared::model::{ConfigDto, HdHomeRunDeviceOverview}; use shared::utils::set_sanitize_sensitive_info; -use crate::model::{macros, ConfigApi, LibraryConfig, ReverseProxyConfig, ScheduleConfig}; -use crate::model::{HdHomeRunConfig, IpCheckConfig, LogConfig, MessagingConfig, ProxyConfig, VideoConfig, WebUiConfig}; -use crate::{utils}; +use std::borrow::Cow; +use std::path::{Path, PathBuf}; const DEFAULT_BACKUP_DIR: &str = "backup"; @@ -128,6 +128,12 @@ impl Config { pub fn is_geoip_enabled(&self) -> bool { self.reverse_proxy.as_ref().is_some_and(|r| r.geoip.as_ref().is_some_and(|g| g.enabled)) } + + pub fn get_disabled_headers(&self) -> Option { + self.reverse_proxy + .as_ref() + .and_then(|r| r.disabled_header.clone()) + } } macros::from_impl!(Config); diff --git a/backend/src/model/xmltv.rs b/backend/src/model/xmltv.rs index ef5103179..182eeef43 100644 --- a/backend/src/model/xmltv.rs +++ b/backend/src/model/xmltv.rs @@ -235,16 +235,12 @@ pub async fn parse_xmltv_for_web_ui_from_url(app_state: &Arc, url: &st headers: HashMap::default(), }; - let disabled_headers = app_state.get_disabled_headers(); - let default_user_agent = app_state.app_config.config.load().default_user_agent.clone(); - match get_remote_content_as_stream( + &app_state.app_config, &client, &input_source, None, &request_url, - disabled_headers.as_ref(), - default_user_agent.as_deref(), ).await { Ok((stream, _url)) => { parse_xmltv_for_web_ui(stream).await diff --git a/backend/src/processing/processor/playlist.rs b/backend/src/processing/processor/playlist.rs index 39b50d409..f5452cc64 100644 --- a/backend/src/processing/processor/playlist.rs +++ b/backend/src/processing/processor/playlist.rs @@ -1,4 +1,4 @@ -use crate::model::{AppConfig, Config, ConfigFavourites, ConfigInput, ConfigRename, TVGuide}; +use crate::model::{AppConfig, Config, ConfigFavourites, ConfigInput, ConfigRename, ReverseProxyDisabledHeaderConfig, TVGuide}; use crate::utils::m3u; use crate::utils::xtream; use crate::utils::{epg, StepMeasureCallback}; @@ -318,7 +318,7 @@ async fn playlist_download_from_input(client: &reqwest::Client, app_config: &Arc let (playlist, errors, persisted) = match input.input_type { InputType::M3u => { - let (p, e) = m3u::download_m3u_playlist(client, config, input).await; + let (p, e) = m3u::download_m3u_playlist(app_config, client, config, input).await; (p, e, false) } InputType::Xtream => xtream::download_xtream_playlist(app_config, client, input, clusters_to_download.as_deref()).await, @@ -458,7 +458,7 @@ async fn download_input_epg(ctx: &PlaylistProcessingContext, input: &Arc, pub event_manager: Option>, pub playlist_state: Option>, + pub disabled_headers: Option, // Coordination processed_inputs: Arc>>>, @@ -863,7 +864,8 @@ async fn process_watch(cfg: &Config, client: &reqwest::Client, target: &ConfigTa pub async fn exec_processing(client: &reqwest::Client, app_config: Arc, targets: Arc, event_manager: Option>, playlist_state: Option>, - update_guard: Option) { + update_guard: Option, + disabled_headers: Option) { let _guard = if let Some(guard) = update_guard { if let Some(permit) = guard.try_playlist() { Some(permit) @@ -887,6 +889,7 @@ pub async fn exec_processing(client: &reqwest::Client, app_config: Arc, client: &reqwest::Client, pli: &PlaylistItem, input: &ConfigInput, errors: &mut Vec, resolve_delay: u16, cluster: XtreamCluster, - default_user_agent: Option<&str>, ) -> 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) { let input_source = InputSource::from(input).with_url(info_url); - result = match xtream::get_xtream_stream_info_content(client, &input_source, true, default_user_agent).await { + result = match xtream::get_xtream_stream_info_content(app_config, client, &input_source, true).await { Ok(content) => Some(content), Err(err) => { errors.push(info_err!("{err}")); diff --git a/backend/src/processing/processor/xtream_series.rs b/backend/src/processing/processor/xtream_series.rs index bd8540b91..9b5c239e0 100644 --- a/backend/src/processing/processor/xtream_series.rs +++ b/backend/src/processing/processor/xtream_series.rs @@ -48,7 +48,6 @@ async fn playlist_resolve_series_info(app_config: &Arc, client: &reqw let mut processed_series_info_count = 0; let mut group_series: HashMap = HashMap::new(); let mut batch = Vec::with_capacity(BATCH_SIZE); - let default_user_agent = app_config.config.load().default_user_agent.clone(); let input = fpl.input; for pli in fpl.items_mut() { @@ -66,13 +65,13 @@ async fn playlist_resolve_series_info(app_config: &Arc, client: &reqw if should_download { processed_series_info_count += 1; if let Some(content) = playlist_resolve_download_playlist_item( + app_config, client, pli, input, errors, resolve_delay, XtreamCluster::Series, - default_user_agent.as_deref(), ) .await { diff --git a/backend/src/processing/processor/xtream_vod.rs b/backend/src/processing/processor/xtream_vod.rs index aaffefb6d..b5dd56c67 100644 --- a/backend/src/processing/processor/xtream_vod.rs +++ b/backend/src/processing/processor/xtream_vod.rs @@ -41,7 +41,6 @@ pub async fn playlist_resolve_vod(app_config: &Arc, let mut last_log_time = Instant::now(); let mut processed_vod_info_count = 0; let mut batch = Vec::with_capacity(BATCH_SIZE); - let default_user_agent = app_config.config.load().default_user_agent.clone(); provider_fpl.source.release_resources(XtreamCluster::Video); @@ -56,13 +55,13 @@ pub async fn playlist_resolve_vod(app_config: &Arc, processed_vod_info_count += 1; if provider_id != 0 { if let Some(content) = playlist_resolve_download_playlist_item( + app_config, client, pli, input, errors, resolve_delay, XtreamCluster::Video, - default_user_agent.as_deref(), ) .await { diff --git a/backend/src/utils/network/epg.rs b/backend/src/utils/network/epg.rs index b66d4efdb..601439b7e 100644 --- a/backend/src/utils/network/epg.rs +++ b/backend/src/utils/network/epg.rs @@ -1,18 +1,17 @@ -use crate::model::TVGuide; use crate::model::{ConfigInput, PersistedEpgSource}; -use crate::utils::{add_prefix_to_filename, prepare_file_path, request}; -use crate::utils::{cleanup_unlisted_files_with_suffix}; -use log::debug; -use shared::error::{TuliproxError, info_err}; -use shared::utils::{sanitize_sensitive_info, short_hash}; -use std::path::PathBuf; -use shared::concat_string; +use crate::model::{TVGuide}; use crate::processing::processor::playlist::PlaylistProcessingContext; use crate::repository::get_input_storage_path; use crate::repository::storage_const; +use crate::utils::{add_prefix_to_filename, prepare_file_path, request}; +use crate::utils::cleanup_unlisted_files_with_suffix; +use log::debug; +use shared::concat_string; +use shared::error::{info_err, TuliproxError}; +use shared::utils::{sanitize_sensitive_info, short_hash}; +use std::path::PathBuf; pub fn get_input_raw_epg_file_path(url: &str, input: &ConfigInput, working_dir: &str) -> std::io::Result { - let file_prefix = short_hash(url); if let Some(persist_path) = input.persist.as_deref() { @@ -28,7 +27,10 @@ pub fn get_input_raw_epg_file_path(url: &str, input: &ConfigInput, working_dir: Ok(download_path.join(format!("{}_{}", file_prefix, storage_const::FILE_EPG))) } -async fn download_epg_file(url: &str, ctx: &PlaylistProcessingContext, input: &ConfigInput, working_dir: &str) -> Result { +async fn download_epg_file(url: &str, ctx: &PlaylistProcessingContext, + input: &ConfigInput, + headers: Option<&reqwest::header::HeaderMap>, + 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).map_err(|e| info_err!("Could not access epg file download directory: {}", e))?; @@ -53,8 +55,7 @@ async fn download_epg_file(url: &str, ctx: &PlaylistProcessingContext, input: &C return Ok(persist_file_path); } debug!("Downloading epg for input '{}'", input.name); - let default_user_agent = ctx.config.config.load().default_user_agent.clone(); - match request::get_input_epg_content_as_file(&ctx.client, input, working_dir, url, &persist_file_path, default_user_agent.as_deref()).await { + match request::get_input_epg_content_as_file(&ctx.config, &ctx.client, input, headers, working_dir, url, &persist_file_path).await { Ok(path) => { ctx.mark_input_downloaded(lock_key.clone()).await; Ok(path) @@ -63,7 +64,9 @@ async fn download_epg_file(url: &str, ctx: &PlaylistProcessingContext, input: &C } } -pub async fn get_xmltv(ctx: &PlaylistProcessingContext, input: &ConfigInput, working_dir: &str) -> (Option, Vec) { +pub async fn get_xmltv(ctx: &PlaylistProcessingContext, input: &ConfigInput, + headers: Option<&reqwest::header::HeaderMap>, + working_dir: &str) -> (Option, Vec) { match &input.epg { None => (None, vec![]), Some(epg_config) => { @@ -72,7 +75,7 @@ pub async fn get_xmltv(ctx: &PlaylistProcessingContext, input: &ConfigInput, wor let mut stored_file_paths = vec![]; for epg_source in &epg_config.sources { - match download_epg_file(&epg_source.url, ctx, input, working_dir).await { + match download_epg_file(&epg_source.url, ctx, input, headers, 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/m3u.rs b/backend/src/utils/network/m3u.rs index 7b013a456..253ce24b2 100644 --- a/backend/src/utils/network/m3u.rs +++ b/backend/src/utils/network/m3u.rs @@ -1,4 +1,4 @@ -use crate::model::{Config, ConfigInput, InputSource}; +use crate::model::{AppConfig, Config, ConfigInput, InputSource}; use crate::processing::parser::m3u; use crate::utils::prepare_file_path; use crate::utils::request; @@ -7,6 +7,7 @@ use shared::model::PlaylistGroup; use std::sync::Arc; pub async fn download_m3u_playlist( + app_config: &Arc, client: &reqwest::Client, cfg: &Arc, input: &ConfigInput, @@ -20,11 +21,11 @@ pub async fn download_m3u_playlist( }; let persist_file_path = prepare_file_path(input.persist.as_deref(), working_dir, ""); match request::get_input_text_content_as_stream( + app_config, client, &input_source, working_dir, persist_file_path, - cfg.default_user_agent.as_deref(), ) .await { diff --git a/backend/src/utils/network/request.rs b/backend/src/utils/network/request.rs index df7b2863f..8b9cc42f7 100644 --- a/backend/src/utils/network/request.rs +++ b/backend/src/utils/network/request.rs @@ -1,14 +1,16 @@ use crate::api::model::persist_pipe_stream::tee_dyn_reader; -use crate::model::ConfigInput; +use crate::api::model::AppState; use crate::model::{format_elapsed_time, AppConfig, InputSource, ReverseProxyDisabledHeaderConfig}; +use crate::model::{ConfigInput, ResourceRetryConfig}; use crate::utils::compression::compression_utils::{is_deflate, is_gzip}; -use crate::utils::{async_file_reader, async_file_writer, debug_if_enabled, IO_BUFFER_SIZE}; +use crate::utils::{async_file_reader, async_file_writer, debug_if_enabled}; use crate::utils::{get_file_path, persist_file}; +use axum::http::header::RETRY_AFTER; use futures::{StreamExt, TryStreamExt}; use log::{debug, error, log_enabled, trace, Level}; use reqwest::header::CONTENT_ENCODING; use reqwest::header::{HeaderMap, HeaderName, HeaderValue}; -use reqwest::StatusCode; +use reqwest::{StatusCode}; use shared::error::{notify_err_res, string_to_io_error, TuliproxError}; use shared::model::{InputFetchMethod, DEFAULT_USER_AGENT}; use shared::utils::{ @@ -74,13 +76,124 @@ pub fn content_type_from_ext(ext: &str) -> &'static str { } } + +#[allow(clippy::cast_possible_truncation, clippy::cast_sign_loss, clippy::cast_precision_loss)] +pub fn calculate_retry_backoff(base_delay_ms: u64, multiplier: f64, attempt: u32) -> u64 { + let base = base_delay_ms.max(1); + if multiplier <= 1.0 { + return base; + } + let delay = (base as f64) * multiplier.powi(i32::try_from(attempt).unwrap_or(i32::MAX)); + if !delay.is_finite() || delay < 1.0 { + base + } else if delay >= u64::MAX as f64 { + u64::MAX + } else { + delay as u64 + } +} + +pub async fn send_with_retry( + app_config: &Arc, + url: &Url, + mut send: impl FnMut() -> reqwest::RequestBuilder, +) -> Result { + let config = app_config.config.load(); + let (max_attempts, backoff_ms, backoff_multiplier) = config + .reverse_proxy + .as_ref() + .map_or_else( + ResourceRetryConfig::get_default_retry_values, + |rp| rp.resource_retry.get_retry_values(), + ); + drop(config); + + for attempt in 0..max_attempts { + match send().send().await { + Ok(response) => { + let status = response.status(); + + if status.is_success() { + return Ok(response); + } + + let should_retry = status.is_server_error() + || matches!( + status, + StatusCode::REQUEST_TIMEOUT + | StatusCode::TOO_EARLY + | StatusCode::TOO_MANY_REQUESTS + ); + + if attempt < max_attempts - 1 && should_retry { + let wait_dur = response + .headers() + .get(RETRY_AFTER) + .and_then(|h| h.to_str().ok()) + .and_then(|s| s.parse::().ok()) + .map_or_else( + || { + let delay = calculate_retry_backoff( + backoff_ms, + backoff_multiplier, + attempt, + ); + Duration::from_millis(delay) + }, + Duration::from_secs, + ); + + tokio::time::sleep(wait_dur).await; + continue; + } + + return Err(string_to_io_error(format!( + "Request failed with status {} {}", + format_http_status(status), + sanitize_sensitive_info(url.as_str()) + ))); + } + + Err(err) => { + if (err.is_timeout() || err.is_connect()) && attempt < max_attempts - 1 { + let delay = calculate_retry_backoff( + backoff_ms, + backoff_multiplier, + attempt, + ); + tokio::time::sleep(Duration::from_millis(delay)).await; + continue; + } + + error!( + "Received failure from server {}: {}", + sanitize_sensitive_info(url.as_str()), + sanitize_sensitive_info(err.to_string().as_str()) + ); + + return Err(string_to_io_error(format!( + "Request failed: {} {}", + sanitize_sensitive_info(url.as_str()), + sanitize_sensitive_info(err.to_string().as_str()) + ))); + } + } + } + + Err(string_to_io_error(format!( + "Failed to download file from {} after all retry attempts", + sanitize_sensitive_info(url.as_str()) + ))) +} + pub async fn get_input_epg_content_as_file( + app_config: &Arc, client: &reqwest::Client, input: &ConfigInput, + headers: Option<&HeaderMap>, working_dir: &str, url_str: &str, persist_filepath: &Path, - default_user_agent: Option<&str>, ) -> Result { debug_if_enabled!( "getting input epg content working_dir: {}, url: {}", @@ -89,13 +202,14 @@ pub async fn get_input_epg_content_as_file( ); if url_str.parse::().is_ok() { match download_epg_content_as_file( + app_config, client, input, + headers, url_str, persist_filepath, - default_user_agent, ) - .await + .await { Ok(content) => Ok(content), Err(e) => { @@ -132,11 +246,11 @@ pub async fn get_input_epg_content_as_file( } pub async fn get_input_text_content( + app_state: &Arc, client: &reqwest::Client, input: &InputSource, working_dir: &str, persist_filepath: Option, - default_user_agent: Option<&str>, ) -> Result { debug_if_enabled!( "getting input text content working_dir: {}, url: {}", @@ -146,15 +260,14 @@ pub async fn get_input_text_content( if input.url.parse::().is_ok() { match download_text_content( + &app_state.app_config, client, - None, input, None, persist_filepath, false, - default_user_agent, ) - .await + .await { Ok((content, _response_url)) => Ok(content), Err(e) => { @@ -195,11 +308,11 @@ pub async fn get_input_text_content( } pub async fn get_input_text_content_as_stream( + app_config: &Arc, client: &reqwest::Client, input: &InputSource, working_dir: &str, persist_filepath: Option, - default_user_agent: Option<&str>, ) -> Result { debug_if_enabled!( "getting input text content working_dir: {}, url: {}", @@ -209,14 +322,12 @@ pub async fn get_input_text_content_as_stream( if input.url.parse::().is_ok() { match download_text_content_as_stream( + app_config, client, - None, input, - None, persist_filepath, - default_user_agent, ) - .await + .await { Ok((content, _response_url)) => Ok(content), Err(e) => { @@ -245,7 +356,7 @@ pub async fn get_input_text_content_as_stream( ); })), ) - .await; + .await; Some(tee) } else { Some(content) @@ -455,138 +566,79 @@ pub async fn get_local_file_content_as_stream( } } -// pub fn get_local_file_content_blocking(file_path: &PathBuf) -> Result { -// match fs::read(file_path) { -// Ok(content) => decode_local_file_bytes(content).await, -// Err(_) => Err(local_file_not_found(file_path)), -// } -// } - -async fn get_remote_content_as_file( +pub async fn get_remote_content_as_file( + app_config: &Arc, client: &reqwest::Client, input: &ConfigInput, - url: &Url, - file_path: &Path, - default_user_agent: Option<&str>, -) -> Result { - let start_time = Instant::now(); - let request = get_client_request( - client, - input.method, - Some(&input.headers), - url, - None, - None, - default_user_agent, - ); - match request.send().await { - Ok(response) => { - if response.status().is_success() { - // Open a file in write mode - let mut writer = async_file_writer(File::create(file_path).await?); - let mut write_counter = 0; - // Stream the response body in chunks - let mut stream = response.bytes_stream(); - while let Some(chunk) = stream.next().await { - match chunk { - Ok(bytes) => { - write_counter += bytes.len(); - writer.write_all(&bytes).await?; - if write_counter >= IO_BUFFER_SIZE { - writer.flush().await?; - write_counter = 0; - } - } - Err(err) => { - let _ = writer.flush().await; - let _ = writer.shutdown().await; - return Err(string_to_io_error(format!("Failed to read chunk: {err}"))); - } - } - } - - writer.flush().await?; - writer.shutdown().await?; - let elapsed = start_time.elapsed().as_secs(); - debug!( - "File downloaded successfully to {}, took:{}", - file_path.display(), - format_elapsed_time(elapsed) - ); - Ok(file_path.to_path_buf()) - } else { - Err(string_to_io_error(format!( - "Request failed with status {} {}", - format_http_status(response.status()), - sanitize_sensitive_info(url.as_str()) - ))) - } - } - Err(err) => Err(string_to_io_error(format!( - "Request failed: {} {err}", - sanitize_sensitive_info(url.as_str()) - ))), - } -} - -pub type DynReader = Pin>; - -#[allow(clippy::implicit_hasher)] -pub async fn get_remote_content_as_stream( - client: &reqwest::Client, - input: &InputSource, headers: Option<&HeaderMap>, url: &Url, - disabled_headers: Option<&ReverseProxyDisabledHeaderConfig>, - default_user_agent: Option<&str>, -) -> Result<(DynReader, String), Error> { + file_path: &Path, +) -> Result { let custom_headers = headers.map(|h| { h.iter() .map(|(k, v)| (k.as_str().to_string(), v.as_bytes().to_vec())) .collect::>() }); - let merged = get_request_headers( - Some(&input.headers), - custom_headers.as_ref(), - disabled_headers, - default_user_agent, - ); - let headers: HashMap = merged - .iter() - .map(|(k, v)| { - ( - k.as_str().to_string(), - String::from_utf8_lossy(v.as_bytes()).to_string(), - ) - }) - .collect(); - let request = get_client_request( - client, - input.method, - Some(&headers), + let config = app_config.config.load(); + let default_user_agent = config.default_user_agent.clone(); + drop(config); + + let response = send_with_retry( + app_config, url, - None, - None, - default_user_agent, - ); - let response = request.send().await.map_err(std::io::Error::other)?; + || { + get_client_request( + client, + input.method, + Some(&input.headers), + url, + custom_headers.as_ref(), + None, + default_user_agent.as_deref(), + ) + }, + ) + .await?; - if !response.status().is_success() { - return Err(string_to_io_error(format!( - "Request failed with status {} {}", - format_http_status(response.status()), - sanitize_sensitive_info(url.as_str()) - ))); + let start_time = Instant::now(); + let mut writer = async_file_writer(File::create(file_path).await?); + + let mut stream = response.bytes_stream(); + while let Some(chunk) = stream.next().await { + let bytes = chunk.map_err(|e| { + string_to_io_error(format!("Failed to read chunk: {e}")) + })?; + writer.write_all(&bytes).await?; } - let response_url = response.url().to_string(); + writer.flush().await?; + writer.shutdown().await?; + + debug!( + "File downloaded successfully to {}, took {}", + file_path.display(), + format_elapsed_time(start_time.elapsed().as_secs()) + ); + + Ok(file_path.to_path_buf()) +} + +pub type DynReader = Pin>; + +async fn build_decoded_stream_reader( + response: reqwest::Response, +) -> Result { let headers = response.headers(); let header_value = headers.get(CONTENT_ENCODING); - let mut encoding = header_value.and_then(|h| h.to_str().ok()).map(ToString::to_string); + let mut encoding = header_value + .and_then(|h| h.to_str().ok()) + .map(ToString::to_string); - let stream_reader = StreamReader::new(response.bytes_stream().map_err(std::io::Error::other)); + let stream_reader = + StreamReader::new(response.bytes_stream().map_err(std::io::Error::other)); let mut buf_reader = async_file_reader(stream_reader); + let peek = buf_reader.fill_buf().await?; if peek.len() >= 2 { @@ -615,27 +667,85 @@ pub async fn get_remote_content_as_stream( Box::pin(buf_reader) }; - Ok((reader, response_url)) + Ok(reader) } -async fn get_remote_content( + +#[allow(clippy::implicit_hasher)] +pub async fn get_remote_content_as_stream( + app_config: &Arc, + client: &reqwest::Client, + input: &InputSource, + headers: Option<&HeaderMap>, + url: &Url, +) -> Result<(DynReader, String), Error> { + let custom_headers = headers.map(|h| { + h.iter() + .map(|(k, v)| (k.as_str().to_string(), v.as_bytes().to_vec())) + .collect::>() + }); + + let config = app_config.config.load(); + let default_user_agent = config.default_user_agent.clone(); + let disabled_headers = config.get_disabled_headers(); + drop(config); + + let merged = get_request_headers( + Some(&input.headers), + custom_headers.as_ref(), + disabled_headers.as_ref(), + default_user_agent.as_deref(), + ); + + let headers: HashMap = merged + .iter() + .map(|(k, v)| { + ( + k.as_str().to_string(), + String::from_utf8_lossy(v.as_bytes()).to_string(), + ) + }) + .collect(); + + let response = send_with_retry( + app_config, + url, + || { + get_client_request( + client, + input.method, + Some(&headers), + url, + None, + None, + default_user_agent.as_deref(), + ) + }, + ) + .await?; + + let response_url = response.url().to_string(); + + let reader = build_decoded_stream_reader(response).await?; + Ok((reader, response_url)) +} + +async fn get_remote_content( + app_config: &Arc, client: &reqwest::Client, input: &InputSource, headers: Option<&HeaderMap>, url: &Url, - disabled_headers: Option<&ReverseProxyDisabledHeaderConfig>, - default_user_agent: Option<&str>, ) -> Result<(String, String), Error> { let (mut stream, response_url) = get_remote_content_as_stream( + app_config, client, input, headers, url, - disabled_headers, - default_user_agent, ) - .await - .map_err(|e| string_to_io_error(format!("Failed to read content: {e}")))?; + .await + .map_err(|e| string_to_io_error(format!("Failed to read content: {e}")))?; let mut content = String::new(); stream .read_to_string(&mut content) @@ -645,11 +755,12 @@ async fn get_remote_content( } async fn download_epg_content_as_file( + app_config: &Arc, client: &reqwest::Client, input: &ConfigInput, + headers: Option<&HeaderMap>, url_str: &str, persist_filepath: &Path, - default_user_agent: Option<&str>, ) -> Result { if let Ok(url) = url_str.parse::() { if url.scheme() == "file" { @@ -672,7 +783,7 @@ async fn download_epg_content_as_file( }, ) } else { - get_remote_content_as_file(client, input, &url, persist_filepath, default_user_agent) + get_remote_content_as_file(app_config, client, input, headers, &url, persist_filepath) .await } } else { @@ -684,13 +795,12 @@ async fn download_epg_content_as_file( } pub async fn download_text_content( + app_config: &Arc, client: &reqwest::Client, - disabled_headers: Option<&ReverseProxyDisabledHeaderConfig>, input: &InputSource, headers: Option<&HeaderMap>, persist_filepath: Option, trace_log: bool, - default_user_agent: Option<&str>, ) -> Result<(String, String), Error> { let start_time = Instant::now(); let result = if let Ok(url) = input.url.parse::() { @@ -706,14 +816,13 @@ pub async fn download_text_content( } } else { get_remote_content( + app_config, client, input, headers, &url, - disabled_headers, - default_user_agent, ) - .await + .await }; match result { Ok((content, response_url)) => { @@ -751,12 +860,10 @@ pub async fn download_text_content( } pub async fn download_text_content_as_stream( + app_config: &Arc, client: &reqwest::Client, - disabled_headers: Option<&ReverseProxyDisabledHeaderConfig>, input: &InputSource, - headers: Option<&HeaderMap>, persist_filepath: Option, - default_user_agent: Option<&str>, ) -> Result<(DynReader, String), Error> { if let Ok(url) = input.url.parse::() { let result = if url.scheme() == "file" { @@ -771,14 +878,13 @@ pub async fn download_text_content_as_stream( } } else { get_remote_content_as_stream( + app_config, client, input, - headers, + None, &url, - disabled_headers, - default_user_agent, ) - .await + .await }; match result { Ok((content, response_url)) => { @@ -790,7 +896,7 @@ pub async fn download_text_content_as_stream( debug!("Persisted {size} bytes"); })), ) - .await; + .await; Ok((tee_reader, response_url)) } else { Ok((content, response_url)) @@ -807,27 +913,25 @@ pub async fn download_text_content_as_stream( } async fn download_json_content( + app_config: &Arc, client: &reqwest::Client, - disabled_headers: Option<&ReverseProxyDisabledHeaderConfig>, input: &InputSource, persist_filepath: Option, trace_log: bool, - default_user_agent: Option<&str>, ) -> Result { debug_if_enabled!( "Downloading json content from {}", sanitize_sensitive_info(&input.url) ); match download_text_content( + app_config, client, - disabled_headers, input, None, persist_filepath, trace_log, - default_user_agent, ) - .await + .await { Ok((content, _response_url)) => match serde_json::from_str::(&content) { Ok(value) => Ok(value), @@ -838,22 +942,20 @@ async fn download_json_content( } pub async fn get_input_json_content( + app_config: &Arc, client: &reqwest::Client, - disabled_headers: Option<&ReverseProxyDisabledHeaderConfig>, input: &InputSource, persist_filepath: Option, trace_log: bool, - default_user_agent: Option<&str>, ) -> Result { match download_json_content( + app_config, client, - disabled_headers, input, persist_filepath, trace_log, - default_user_agent, ) - .await + .await { Ok(content) => Ok(content), Err(e) => notify_err_res!( @@ -865,25 +967,22 @@ pub async fn get_input_json_content( } async fn download_json_content_as_stream( + app_config: &Arc, client: &reqwest::Client, - disabled_headers: Option<&ReverseProxyDisabledHeaderConfig>, input: &InputSource, persist_filepath: Option, - default_user_agent: Option<&str>, ) -> Result { debug_if_enabled!( "Downloading json content as stream from {}", sanitize_sensitive_info(&input.url) ); match download_text_content_as_stream( + app_config, client, - disabled_headers, input, - None, persist_filepath, - default_user_agent, ) - .await + .await { Ok((reader, _response_url)) => Ok(reader), Err(err) => Err(err), @@ -891,20 +990,18 @@ async fn download_json_content_as_stream( } pub async fn get_input_json_content_as_stream( + app_config: &Arc, client: &reqwest::Client, - disabled_headers: Option<&ReverseProxyDisabledHeaderConfig>, input: &InputSource, persist_filepath: Option, - default_user_agent: Option<&str>, ) -> Result { match download_json_content_as_stream( + app_config, client, - disabled_headers, input, persist_filepath, - default_user_agent, ) - .await + .await { Ok(stream) => Ok(stream), Err(e) => notify_err_res!( diff --git a/backend/src/utils/network/xtream.rs b/backend/src/utils/network/xtream.rs index d83917877..999237d52 100644 --- a/backend/src/utils/network/xtream.rs +++ b/backend/src/utils/network/xtream.rs @@ -60,8 +60,8 @@ pub fn get_xtream_player_api_info_url(input: &ConfigInput, cluster: XtreamCluste } -pub async fn get_xtream_stream_info_content(client: &reqwest::Client, input: &InputSource, trace_log: bool, default_user_agent: Option<&str>,) -> Result { - match request::download_text_content(client, None, input, None, None, trace_log, default_user_agent).await { +pub async fn get_xtream_stream_info_content(app_config: &Arc, client: &reqwest::Client, input: &InputSource, trace_log: bool) -> Result { + match request::download_text_content(app_config, client, input, None, None, trace_log).await { Ok((content, _response_url)) => Ok(content), Err(err) => Err(err) } @@ -87,7 +87,7 @@ pub async fn get_xtream_stream_info(client: &reqwest::Client, } let input_source = InputSource::from(input).with_url(info_url.to_owned()); - if let Ok(content) = get_xtream_stream_info_content(client, &input_source, false, app_config.config.load().default_user_agent.as_deref()).await { + if let Ok(content) = get_xtream_stream_info_content(app_config, client, &input_source, false).await { if content.is_empty() { return Err(info_err!("Provider returned no response for stream with id: {}/{}/{}", target.name.replace(' ', "_").as_str(), &cluster, pli.get_virtual_id())); @@ -266,13 +266,13 @@ 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: &reqwest::Client, input: &InputSource, username: &str) -> Result, TuliproxError> { - let content = if let Ok(content) = request::get_input_json_content(client, None, input, None, false, cfg.default_user_agent.as_deref()).await { +async fn xtream_login(app_config: &Arc, client: &reqwest::Client, input: &InputSource, username: &str) -> Result, TuliproxError> { + let content = if let Ok(content) = request::get_input_json_content(app_config, client, input, None, false).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(client, None, &input_source_account_info, None, false, cfg.default_user_agent.as_deref()).await { + match request::get_input_json_content(app_config, client, &input_source_account_info, None, false).await { Ok(content) => content, Err(err) => { warn!("Failed to login xtream account {username} {err}"); @@ -286,6 +286,8 @@ async fn xtream_login(cfg: &Config, client: &reqwest::Client, input: &InputSourc exp_date: None, }; + let cfg = app_config.config.load(); + if let Some(user_info) = content.get("user_info") { if let Some(status_value) = user_info.get("status") { if let Some(status) = get_string_from_serde_value(status_value) { @@ -303,7 +305,7 @@ async fn xtream_login(cfg: &Config, client: &reqwest::Client, input: &InputSourc if let Some(exp_value) = user_info.get("exp_date") { if let Some(expiration_timestamp) = get_i64_from_serde_value(exp_value) { login_info.exp_date = Some(expiration_timestamp); - notify_account_expire(login_info.exp_date, cfg, client, username, &input.name).await; + notify_account_expire(login_info.exp_date, &cfg, client, username, &input.name).await; } } } @@ -356,7 +358,7 @@ pub async fn download_xtream_playlist(app_config: &Arc, client: &reqw check_alias_user_state(&cfg, client, input).await; - if let Err(err) = xtream_login(&cfg, client, &input_source_login, username).await { + if let Err(err) = xtream_login(app_config, client, &input_source_login, username).await { error!("Could not log in with xtream user {username} for provider {}. {err}", input.name); return (Vec::with_capacity(0), vec![err], false); } @@ -378,8 +380,8 @@ pub async fn download_xtream_playlist(app_config: &Arc, client: &reqw working_dir, format!("{stream}_").as_str()); match futures::join!( - request::get_input_json_content_as_stream(client, None, &input_source_category, category_file_path, cfg.default_user_agent.as_deref()), - request::get_input_json_content_as_stream(client, None, &input_source_stream, stream_file_path, cfg.default_user_agent.as_deref()) + request::get_input_json_content_as_stream(app_config, client, &input_source_category, category_file_path), + request::get_input_json_content_as_stream(app_config, client, &input_source_stream, stream_file_path) ) { (Ok(category_content), Ok(stream_content)) => { if cfg.disk_based_processing { diff --git a/shared/src/error/tuliprox_error.rs b/shared/src/error/tuliprox_error.rs index 8af67b703..3ab7af9b1 100644 --- a/shared/src/error/tuliprox_error.rs +++ b/shared/src/error/tuliprox_error.rs @@ -1,5 +1,6 @@ use std::error::Error; use std::fmt::{Display, Formatter, Result}; +use crate::utils::sanitize_sensitive_info; #[macro_export] macro_rules! get_errors_notify_message { @@ -94,10 +95,8 @@ macro_rules! handle_tuliprox_error_result { } } } - pub use handle_tuliprox_error_result; - #[derive(Debug, Copy, Clone, PartialEq, Eq)] pub enum TuliproxErrorKind { // do not send with messaging @@ -128,12 +127,12 @@ impl Error for TuliproxError {} pub fn to_io_error(err: E) -> std::io::Error where E: std::error::Error, -{ std::io::Error::other(err.to_string()) } +{ std::io::Error::other(sanitize_sensitive_info(&err.to_string())) } pub fn str_to_io_error(err: &str) -> std::io::Error { - std::io::Error::other(err.to_string()) + std::io::Error::other(sanitize_sensitive_info(err)) } pub fn string_to_io_error(err: String) -> std::io::Error { - std::io::Error::other(err) + std::io::Error::other(sanitize_sensitive_info(&err)) } diff --git a/shared/src/model/short_epg.rs b/shared/src/model/short_epg.rs index e19d71dc9..96ad2859d 100644 --- a/shared/src/model/short_epg.rs +++ b/shared/src/model/short_epg.rs @@ -17,14 +17,9 @@ pub struct ShortEpgDto { pub stream_id: String, } -#[derive(Debug, serde::Serialize, serde::Deserialize, Default)] -pub struct ShortEpgListingsDto { - pub epg_listings : Vec, -} - #[derive(Debug, serde::Serialize, serde::Deserialize, Default)] pub struct ShortEpgResultDto { - pub data: ShortEpgListingsDto, + pub epg_listings : Vec, #[serde(default, serialize_with = "serialize_option_string_as_null_if_empty")] pub error: Option, } @@ -32,9 +27,7 @@ pub struct ShortEpgResultDto { impl ShortEpgResultDto { pub fn new(epg_listings: Vec) -> Self { Self { - data: ShortEpgListingsDto { - epg_listings, - }, + epg_listings, error: None, } }