From dc8ef8f8c6fbd1e9b7077142db1d44c4ddeb2b99 Mon Sep 17 00:00:00 2001 From: euzu Date: Wed, 30 Apr 2025 13:32:55 +0200 Subject: [PATCH 1/2] proxy settings --- CHANGELOG.md | 7 ++++++ src/api/endpoints/download_api.rs | 10 +++++--- src/api/main_api.rs | 3 ++- src/main.rs | 29 ++++++++++++---------- src/messaging.rs | 25 ++++++++++--------- src/model/config.rs | 37 ++++++++++++++++++++++++++++ src/processing/playlist_watch.rs | 9 ++++--- src/processing/processor/playlist.rs | 12 ++++----- src/utils/network/request.rs | 27 +++++++++++++++++++- 9 files changed, 118 insertions(+), 41 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 535402116..a5e8adc6b 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -99,6 +99,13 @@ hdhomerun: If `caption` is provided, its value is read from `title` if available, otherwise from `name`. When setting `caption`, both `title` and `name` are updated.” - Counter has now an attribute padding. Which fills the number like 001. +- Added proxy configuration for all outgoing requests in `config.yml`. supported http, https, socks5 proxies. +```yaml +proxy: + url: socks5://192.168.1.6:8123 + username: uname # <- optional basic auth + password: secret # <- optional basic auth +``` # 2.2.5 (2025-03-27) - fixed web ui playlist regexp search diff --git a/src/api/endpoints/download_api.rs b/src/api/endpoints/download_api.rs index 54c04e850..e724aa576 100644 --- a/src/api/endpoints/download_api.rs +++ b/src/api/endpoints/download_api.rs @@ -1,6 +1,6 @@ use crate::api::model::app_state::AppState; use crate::api::model::download::{DownloadQueue, FileDownload, FileDownloadRequest}; -use crate::model::config::VideoDownloadConfig; +use crate::model::config::{ConfigProxy, VideoDownloadConfig}; use crate::utils::network::request; use tokio::sync::RwLock; use futures::stream::TryStreamExt; @@ -13,6 +13,7 @@ use std::sync::Arc; use std::{fs}; use axum::response::IntoResponse; use crate::m3u_filter_error::to_io_error; +use crate::utils::network::request::create_client; async fn download_file(active: Arc>>, client: &reqwest::Client) -> Result<(), String> { let file_download = { active.read().await.as_ref().unwrap().clone() }; @@ -61,13 +62,14 @@ async fn download_file(active: Arc>>, client: &reqwe } } -async fn run_download_queue(download_cfg: &VideoDownloadConfig, download_queue: &Arc) -> Result<(), String> { +async fn run_download_queue(proxy_config: Option<&ConfigProxy>, download_cfg: &VideoDownloadConfig, download_queue: &Arc) -> Result<(), String> { let next_download = download_queue.as_ref().queue.lock().await.pop_front(); if next_download.is_some() { { *download_queue.as_ref().active.write().await = next_download; } let headers = request::get_request_headers(Some(&download_cfg.headers), None); let dq = Arc::clone(download_queue); - match reqwest::Client::builder().default_headers(headers).build() { + + match create_client(proxy_config).default_headers(headers).build() { Ok(client) => { tokio::spawn(async move { loop { @@ -121,7 +123,7 @@ pub async fn queue_download_file( Some(file_download) => { app_state.downloads.queue.lock().await.push_back(file_download.clone()); if app_state.downloads.active.read().await.is_none() { - match run_download_queue(download_cfg, &app_state.downloads).await { + match run_download_queue(app_state.config.proxy.as_ref(), download_cfg, &app_state.downloads).await { Ok(()) => {} Err(err) => return (axum::http::StatusCode::INTERNAL_SERVER_ERROR, axum::Json(json!({"error": err}))).into_response(), } diff --git a/src/api/main_api.rs b/src/api/main_api.rs index 7719f6687..a5316935d 100644 --- a/src/api/main_api.rs +++ b/src/api/main_api.rs @@ -27,6 +27,7 @@ use axum::Router; use tokio::sync::Mutex; use tower_governor::key_extractor::SmartIpKeyExtractor; use crate::api::api_utils::{get_build_time, get_server_time}; +use crate::utils::network::request::create_client; use crate::VERSION; fn get_web_dir_path(web_ui_enabled: bool, web_root: &str) -> Result { @@ -71,7 +72,7 @@ async fn create_shared_data(cfg: &Arc) -> AppState { let active_users = Arc::new(ActiveUserManager::new(cfg)); let active_provider = Arc::new(ActiveProviderManager::new(cfg).await); - let mut builder = Client::builder().http1_only(); + let mut builder = create_client(cfg.proxy.as_ref()).http1_only(); if cfg.connect_timeout_secs > 0 { builder = builder.connect_timeout(Duration::from_secs(u64::from(cfg.connect_timeout_secs))); } diff --git a/src/main.rs b/src/main.rs index 08027ecfd..13dbd11d4 100644 --- a/src/main.rs +++ b/src/main.rs @@ -1,4 +1,4 @@ -#![warn(clippy::all,clippy::pedantic)] +#![warn(clippy::all, clippy::pedantic)] #![allow(clippy::module_name_repetitions)] #![allow(clippy::must_use_candidate)] #![allow(clippy::return_self_not_must_use)] @@ -9,21 +9,21 @@ mod modules; include_modules!(); -use std::fs::File; -use std::path::{Path, PathBuf}; -use std::sync::Arc; -use chrono::{DateTime, Utc}; use crate::auth::password::generate_password; use crate::model::config::{validate_targets, Config, HealthcheckConfig, LogLevelConfig, ProcessTargets}; use crate::model::healthcheck::Healthcheck; use crate::processing::processor::playlist; -use utils::file::config_reader; +use crate::utils::file::config_reader::config_file_reader; use crate::utils::file::file_utils; -use crate::utils::network::request::set_sanitize_sensitive_info; +use crate::utils::network::request::{create_client, set_sanitize_sensitive_info}; +use chrono::{DateTime, Utc}; use clap::Parser; use env_logger::Builder; use log::{error, info, LevelFilter}; -use crate::utils::file::config_reader::config_file_reader; +use std::fs::File; +use std::path::{Path, PathBuf}; +use std::sync::Arc; +use utils::file::config_reader; const LOG_ERROR_LEVEL_MOD: &[&str] = &[ "reqwest::async_impl::client", @@ -80,7 +80,7 @@ struct Args { const VERSION: &str = env!("CARGO_PKG_VERSION"); -const BUILD_TIMESTAMP:&str = env!("VERGEN_BUILD_TIMESTAMP"); +const BUILD_TIMESTAMP: &str = env!("VERGEN_BUILD_TIMESTAMP"); // #[cfg(not(target_env = "msvc"))] // #[global_allocator] @@ -150,13 +150,13 @@ fn main() { // info!("Channel unavailable video loaded from {:?}", cfg.channel_unavailable_file.as_ref().map_or("?", |v| v.as_str())); // } - let rt = tokio::runtime::Runtime::new().unwrap(); + let rt = tokio::runtime::Runtime::new().unwrap(); let () = rt.block_on(async { if args.server { match config_reader::read_api_proxy_config(args.api_proxy, &mut cfg).await { Ok(Some(api_proxy_file)) => { info!("Api Proxy File: {api_proxy_file:?}"); - }, + } Ok(None) => {} Err(err) => exit!("{err}"), } @@ -197,8 +197,11 @@ fn create_directories(cfg: &Config, temp_path: &Path) { } async fn start_in_cli_mode(cfg: Arc, targets: Arc) { - let client = Arc::new(reqwest::Client::new()); - playlist::exec_processing(client, cfg, targets).await; + let client = create_client(cfg.proxy.as_ref()).build().unwrap_or_else(|err| { + error!("Failed to build cient {err}"); + reqwest::Client::new() + }); + playlist::exec_processing(Arc::new(client), cfg, targets).await; } async fn start_in_server_mode(cfg: Arc, targets: Arc) { diff --git a/src/messaging.rs b/src/messaging.rs index ff76df17f..d6bb05448 100644 --- a/src/messaging.rs +++ b/src/messaging.rs @@ -1,6 +1,7 @@ -use crate::model::config::MessagingConfig; +use std::sync::Arc; +use crate::model::config::{MessagingConfig}; use log::{debug, error}; -use reqwest::header; +use reqwest::{header}; #[derive(Debug, Clone, serde::Serialize, serde::Deserialize, PartialEq, Eq)] pub enum MsgKind { @@ -18,13 +19,13 @@ fn is_enabled(kind: &MsgKind, cfg: &MessagingConfig) -> bool { cfg.notify_on.contains(kind) } -fn send_http_post_request(msg: &str, messaging: &MessagingConfig) { +fn send_http_post_request(client: &Arc, msg: &str, messaging: &MessagingConfig) { if let Some(rest) = &messaging.rest { let url = rest.url.clone(); let data = msg.to_owned(); + let the_client = Arc::clone(client); tokio::spawn(async move { - let client = reqwest::Client::new(); - match client + match the_client .post(&url) .header(header::CONTENT_TYPE, mime::APPLICATION_JSON.to_string()) .body(data) @@ -39,6 +40,7 @@ fn send_http_post_request(msg: &str, messaging: &MessagingConfig) { } fn send_telegram_message(msg: &str, messaging: &MessagingConfig) { + // TODO use proxy settings if let Some(telegram) = &messaging.telegram { for chat_id in &telegram.chat_ids { let bot = rustelebot::create_instance(&telegram.bot_token, chat_id); @@ -50,7 +52,7 @@ fn send_telegram_message(msg: &str, messaging: &MessagingConfig) { } } -fn send_pushover_message(msg: &str, messaging: &MessagingConfig) { +fn send_pushover_message(client: &Arc, msg: &str, messaging: &MessagingConfig) { if let Some(pushover) = &messaging.pushover { let url = pushover.url.as_deref().unwrap_or("https://api.pushover.net/1/messages.json").to_string(); let encoded_message: String = url::form_urlencoded::Serializer::new(String::new()) @@ -58,10 +60,9 @@ fn send_pushover_message(msg: &str, messaging: &MessagingConfig) { .append_pair("user", pushover.user.as_str()) .append_pair("message", msg) .finish(); - + let the_client = Arc::clone(client); tokio::spawn(async move { - let client = reqwest::Client::new(); - match client + match the_client .post(url) .header(header::CONTENT_TYPE, mime::APPLICATION_WWW_FORM_URLENCODED.to_string()) .body(encoded_message) @@ -81,12 +82,12 @@ fn send_pushover_message(msg: &str, messaging: &MessagingConfig) { } } -pub fn send_message(kind: &MsgKind, cfg: Option<&MessagingConfig>, msg: &str) { +pub fn send_message(client: &Arc, kind: &MsgKind, cfg: Option<&MessagingConfig>, msg: &str) { if let Some(messaging) = cfg { if is_enabled(kind, messaging) { send_telegram_message(msg, messaging); - send_http_post_request(msg, messaging); - send_pushover_message(msg, messaging); + send_http_post_request(client, msg, messaging); + send_pushover_message(client, msg, messaging); } } } diff --git a/src/model/config.rs b/src/model/config.rs index f3ad742a9..03e4de82e 100644 --- a/src/model/config.rs +++ b/src/model/config.rs @@ -1765,6 +1765,38 @@ impl WebUiConfig { } } +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, Default)] +#[serde(deny_unknown_fields)] +pub struct ConfigProxy { + pub url: String, + pub username: Option, + pub password: Option, +} + +impl ConfigProxy { + fn prepare(&mut self) -> Result<(), M3uFilterError> { + if self.username.is_some() || self.password.is_some() { + if let (Some(username), Some(password)) = (self.username.as_ref(), self.password.as_ref()) { + let uname = username.trim(); + let pwd = password.trim(); + if uname.is_empty() || pwd.is_empty() { + return Err(M3uFilterError::new(M3uFilterErrorKind::Info,"Proxy credentials missing".to_string())); + } + self.username = Some(uname.to_string()); + self.password = Some(pwd.to_string()); + } else { + return Err(M3uFilterError::new(M3uFilterErrorKind::Info,"Proxy credentials missing".to_string())); + } + } + + self.url = self.url.trim().to_string(); + if self.url.is_empty() { + return Err(M3uFilterError::new(M3uFilterErrorKind::Info,"Proxy url missing".to_string())); + } + Ok(()) + } +} + #[derive(Debug, Clone, serde::Serialize, serde::Deserialize, Default)] #[serde(deny_unknown_fields)] pub struct Config { @@ -1801,6 +1833,8 @@ pub struct Config { pub reverse_proxy: Option, #[serde(default, skip_serializing_if = "Option::is_none")] pub hdhomerun: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub proxy: Option, #[serde(skip)] pub t_api_proxy: Arc>>, #[serde(skip)] @@ -2074,6 +2108,9 @@ impl Config { if let Some(reverse_proxy) = self.reverse_proxy.as_mut() { reverse_proxy.prepare(&self.working_dir)?; } + if let Some(proxy) = &mut self.proxy { + proxy.prepare()?; + } self.prepare_hdhomerun()?; self.api.prepare(); self.prepare_api_web_root(); diff --git a/src/processing/playlist_watch.rs b/src/processing/playlist_watch.rs index 162b68043..5f0c30bac 100644 --- a/src/processing/playlist_watch.rs +++ b/src/processing/playlist_watch.rs @@ -1,5 +1,6 @@ use std::collections::BTreeSet; use std::path::{Path}; +use std::sync::Arc; use log::{error, info}; use crate::messaging::{MsgKind, send_message}; use crate::model::config::Config; @@ -8,7 +9,7 @@ use crate::utils::bincode_utils::{bincode_deserialize, bincode_serialize}; use crate::utils::file::file_utils; use crate::utils::file::file_utils::sanitize_filename; -pub fn process_group_watch(cfg: &Config, target_name: &str, pl: &PlaylistGroup) { +pub fn process_group_watch(client: &Arc, cfg: &Config, target_name: &str, pl: &PlaylistGroup) { let mut new_tree = BTreeSet::new(); pl.channels.iter().for_each(|chan| { let header = &chan.header; @@ -28,7 +29,7 @@ pub fn process_group_watch(cfg: &Config, target_name: &str, pl: &PlaylistGroup) let removed_difference: BTreeSet = loaded_tree.difference(&new_tree).cloned().collect(); if !added_difference.is_empty() || !removed_difference.is_empty() { changed = true; - handle_watch_notification(cfg, &added_difference, &removed_difference, target_name, &pl.title); + handle_watch_notification(client, cfg, &added_difference, &removed_difference, target_name, &pl.title); } } else { error!("failed to load watch_file {}", &path.to_str().unwrap_or_default()); @@ -52,7 +53,7 @@ pub fn process_group_watch(cfg: &Config, target_name: &str, pl: &PlaylistGroup) } } -fn handle_watch_notification(cfg: &Config, added: &BTreeSet, removed: &BTreeSet, target_name: &str, group_name: &str) { +fn handle_watch_notification(client: &Arc, cfg: &Config, added: &BTreeSet, removed: &BTreeSet, target_name: &str, group_name: &str) { let added_entries = added.iter().map(std::string::ToString::to_string).collect::>().join("\n\t"); let removed_entries = removed.iter().map(std::string::ToString::to_string).collect::>().join("\n\t"); @@ -71,7 +72,7 @@ fn handle_watch_notification(cfg: &Config, added: &BTreeSet, removed: &B if !message.is_empty() { let msg = format!("Changes {}/{}\n{}", target_name, group_name, message.join("")); info!("{}", &msg); - send_message(&MsgKind::Watch, cfg.messaging.as_ref(), &msg); + send_message(client, &MsgKind::Watch, cfg.messaging.as_ref(), &msg); } } diff --git a/src/processing/processor/playlist.rs b/src/processing/processor/playlist.rs index 66583f8af..78f658526 100644 --- a/src/processing/processor/playlist.rs +++ b/src/processing/processor/playlist.rs @@ -563,7 +563,7 @@ async fn process_playlist_for_target(client: Arc, step.tick("Assigned channel counter"); map_playlist_counter(target, &mut flat_new_playlist); step.tick("Processed group watches"); - process_watch(target, cfg, &flat_new_playlist); + process_watch(&client, target, cfg, &flat_new_playlist); step.tick("Persisting playlists"); let result = persist_playlist(&mut flat_new_playlist, flatten_tvguide(&new_epg).as_ref(), target, cfg).await; step.stop(); @@ -585,7 +585,7 @@ fn process_epg(processed_fetched_playlists: &mut Vec) -> (Vec) { +fn process_watch(client: &Arc, target: &ConfigTarget, cfg: &Config, new_playlist: &Vec) { if target.t_watch_re.is_some() { if default_as_default().eq_ignore_ascii_case(&target.name) { error!("cant watch a target with no unique name"); @@ -593,7 +593,7 @@ fn process_watch(target: &ConfigTarget, cfg: &Config, new_playlist: &Vec, cfg: Arc, targets: Arc) { let start_time = Instant::now(); - let (stats, errors) = process_sources(client, cfg.clone(), targets.clone()).await; + let (stats, errors) = process_sources(Arc::clone(&client), cfg.clone(), targets.clone()).await; // log errors for err in &errors { error!("{}", err.message); @@ -611,12 +611,12 @@ pub async fn exec_processing(client: Arc, cfg: Arc, tar // print stats info!("{stats_msg}"); // send stats - send_message(&MsgKind::Stats, cfg.messaging.as_ref(), stats_msg.as_str()); + send_message(&client, &MsgKind::Stats, cfg.messaging.as_ref(), stats_msg.as_str()); } // send errors if let Some(message) = get_errors_notify_message!(errors, 255) { if let Ok(error_msg) = serde_json::to_string(&serde_json::Value::Object(serde_json::map::Map::from_iter([("errors".to_string(), serde_json::Value::String(message))]))) { - send_message(&MsgKind::Error, cfg.messaging.as_ref(), error_msg.as_str()); + send_message(&client, &MsgKind::Error, cfg.messaging.as_ref(), error_msg.as_str()); } } let elapsed = start_time.elapsed().as_secs(); diff --git a/src/utils/network/request.rs b/src/utils/network/request.rs index 7ecbb663e..d590c892e 100644 --- a/src/utils/network/request.rs +++ b/src/utils/network/request.rs @@ -16,7 +16,7 @@ use url::Url; use crate::m3u_filter_error::create_m3u_filter_error_result; use crate::m3u_filter_error::{str_to_io_error, M3uFilterError, M3uFilterErrorKind}; -use crate::model::config::{ConfigInput, InputFetchMethod}; +use crate::model::config::{ConfigInput, ConfigProxy, InputFetchMethod}; use crate::model::stats::format_elapsed_time; use crate::repository::storage::{get_input_storage_path, short_hash}; use crate::repository::storage_const; @@ -496,6 +496,31 @@ pub fn get_base_url_from_str(url: &str) -> Option { } } +pub fn create_client(proxy_config: Option<&ConfigProxy>) -> reqwest::ClientBuilder { + let client = reqwest::Client::builder(); + if let Some(proxy_cfg) = proxy_config { + let proxy = match reqwest::Proxy::all(&proxy_cfg.url) { + Ok(proxy) => { + if let (Some(username), Some(password)) = (&proxy_cfg.username, &proxy_cfg.password) { + Some(proxy.basic_auth(username, password)) + } else { + Some(proxy) + } + } + Err(err) => { + error!("Failed to create proxy {}, {err}", &proxy_cfg.url); + None + } + }; + return if let Some(prxy) = proxy { + client.proxy(prxy) + } else { + client + }; + } + client +} + #[cfg(test)] mod tests { use crate::utils::network::request::{get_base_url_from_str, replace_url_extension, sanitize_sensitive_info}; From 86ecee997cfb24027a74ac9d83035e38c87fcec6 Mon Sep 17 00:00:00 2001 From: euzu Date: Wed, 30 Apr 2025 13:40:26 +0200 Subject: [PATCH 2/2] proxy settings --- src/main.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/main.rs b/src/main.rs index 13dbd11d4..8d8b7d7ba 100644 --- a/src/main.rs +++ b/src/main.rs @@ -198,7 +198,7 @@ fn create_directories(cfg: &Config, temp_path: &Path) { async fn start_in_cli_mode(cfg: Arc, targets: Arc) { let client = create_client(cfg.proxy.as_ref()).build().unwrap_or_else(|err| { - error!("Failed to build cient {err}"); + error!("Failed to build client {err}"); reqwest::Client::new() }); playlist::exec_processing(Arc::new(client), cfg, targets).await;