mirror of
https://github.com/euzu/tuliprox.git
synced 2026-09-29 04:22:33 +02:00
Merge pull request #235 from euzu/featue/proxy_connection
proxy settings
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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<RwLock<Option<FileDownload>>>, 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<RwLock<Option<FileDownload>>>, client: &reqwe
|
||||
}
|
||||
}
|
||||
|
||||
async fn run_download_queue(download_cfg: &VideoDownloadConfig, download_queue: &Arc<DownloadQueue>) -> Result<(), String> {
|
||||
async fn run_download_queue(proxy_config: Option<&ConfigProxy>, download_cfg: &VideoDownloadConfig, download_queue: &Arc<DownloadQueue>) -> 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(),
|
||||
}
|
||||
|
||||
+2
-1
@@ -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<PathBuf, std::io::Error> {
|
||||
@@ -71,7 +72,7 @@ async fn create_shared_data(cfg: &Arc<Config>) -> 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)));
|
||||
}
|
||||
|
||||
+16
-13
@@ -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<Config>, targets: Arc<ProcessTargets>) {
|
||||
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 client {err}");
|
||||
reqwest::Client::new()
|
||||
});
|
||||
playlist::exec_processing(Arc::new(client), cfg, targets).await;
|
||||
}
|
||||
|
||||
async fn start_in_server_mode(cfg: Arc<Config>, targets: Arc<ProcessTargets>) {
|
||||
|
||||
+13
-12
@@ -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<reqwest::Client>, 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<reqwest::Client>, 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<reqwest::Client>, 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);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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<String>,
|
||||
pub password: Option<String>,
|
||||
}
|
||||
|
||||
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<ReverseProxyConfig>,
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub hdhomerun: Option<HdHomeRunConfig>,
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub proxy: Option<ConfigProxy>,
|
||||
#[serde(skip)]
|
||||
pub t_api_proxy: Arc<RwLock<Option<ApiProxyConfig>>>,
|
||||
#[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();
|
||||
|
||||
@@ -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<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;
|
||||
@@ -28,7 +29,7 @@ pub fn process_group_watch(cfg: &Config, target_name: &str, pl: &PlaylistGroup)
|
||||
let removed_difference: BTreeSet<String> = 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<String>, removed: &BTreeSet<String>, target_name: &str, group_name: &str) {
|
||||
fn handle_watch_notification(client: &Arc<reqwest::Client>, cfg: &Config, added: &BTreeSet<String>, removed: &BTreeSet<String>, target_name: &str, group_name: &str) {
|
||||
let added_entries = added.iter().map(std::string::ToString::to_string).collect::<Vec<String>>().join("\n\t");
|
||||
let removed_entries = removed.iter().map(std::string::ToString::to_string).collect::<Vec<String>>().join("\n\t");
|
||||
|
||||
@@ -71,7 +72,7 @@ fn handle_watch_notification(cfg: &Config, added: &BTreeSet<String>, 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);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -563,7 +563,7 @@ async fn process_playlist_for_target(client: Arc<reqwest::Client>,
|
||||
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<FetchedPlaylist>) -> (Vec<E
|
||||
(new_epg, new_playlist)
|
||||
}
|
||||
|
||||
fn process_watch(target: &ConfigTarget, cfg: &Config, new_playlist: &Vec<PlaylistGroup>) {
|
||||
fn process_watch(client: &Arc<reqwest::Client>, target: &ConfigTarget, cfg: &Config, new_playlist: &Vec<PlaylistGroup>) {
|
||||
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<Playlis
|
||||
let watch_re = target.t_watch_re.as_ref().unwrap();
|
||||
for pl in new_playlist {
|
||||
if watch_re.iter().any(|r| r.is_match(&pl.title)) {
|
||||
process_group_watch(cfg, &target.name, pl);
|
||||
process_group_watch(client, cfg, &target.name, pl);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -602,7 +602,7 @@ fn process_watch(target: &ConfigTarget, cfg: &Config, new_playlist: &Vec<Playlis
|
||||
|
||||
pub async fn exec_processing(client: Arc<reqwest::Client>, cfg: Arc<Config>, targets: Arc<ProcessTargets>) {
|
||||
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<reqwest::Client>, cfg: Arc<Config>, 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();
|
||||
|
||||
@@ -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<String> {
|
||||
}
|
||||
}
|
||||
|
||||
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};
|
||||
|
||||
Reference in New Issue
Block a user