diff --git a/backend/src/api/config_watch.rs b/backend/src/api/config_watch.rs index 04bfd57a8..d5133b763 100644 --- a/backend/src/api/config_watch.rs +++ b/backend/src/api/config_watch.rs @@ -1,5 +1,5 @@ -use crate::api::model::app_state::AppState; -use crate::utils; +use crate::utils::exit; +use crate::api::model::app_state::{update_app_state, AppState}; use crate::utils::{is_directory, read_config_file, read_sources_file}; use log::{debug, error, info}; use notify::event::{AccessKind, AccessMode}; @@ -9,6 +9,7 @@ use std::collections::HashMap; use std::path::{Path, PathBuf}; use std::sync::{mpsc, Arc}; use crate::model::{Config, SourcesConfig}; +use crate::utils; enum ConfigFile { Config, @@ -65,7 +66,7 @@ impl ConfigFile { let mut config: Config = Config::from(config_dto); config.prepare(paths.config_path.as_str())?; info!("Loaded config file {config_file}"); - app_state.set_config(config).await?; + update_app_state(app_state, config).await?; Ok(()) } @@ -75,7 +76,11 @@ impl ConfigFile { let sources_dto = read_sources_file(sources_file, true, true)?; let sources: SourcesConfig = sources_dto.into(); info!("Loaded sources file {sources_file}"); + + let targets = sources.validate_targets(Some(&app_state.forced_targets.load().target_names)).unwrap_or_else(|err| exit!("{}", err)); + app_state.forced_targets.store(Arc::new(targets)); app_state.app_config.set_sources(sources)?; + Ok(()) } diff --git a/backend/src/api/endpoints/hls_api.rs b/backend/src/api/endpoints/hls_api.rs index 6a62ce77d..ebc948db8 100644 --- a/backend/src/api/endpoints/hls_api.rs +++ b/backend/src/api/endpoints/hls_api.rs @@ -98,7 +98,7 @@ async fn hls_api_stream( ) -> impl axum::response::IntoResponse + Send { let (user, target) = try_option_bad_request!( app_state.app_config.get_target_for_user(¶ms.username, ¶ms.password), false, - format!("Could not find any user {}", params.username)); + format!("Could not find any user for hls stream {}", params.username)); if user.permission_denied(&app_state) { return create_custom_video_stream_response(&app_state.app_config, CustomVideoStreamType::UserAccountExpired).into_response(); } diff --git a/backend/src/api/endpoints/m3u_api.rs b/backend/src/api/endpoints/m3u_api.rs index 0a6bb6bbc..5ac7dd85a 100644 --- a/backend/src/api/endpoints/m3u_api.rs +++ b/backend/src/api/endpoints/m3u_api.rs @@ -67,7 +67,7 @@ async fn m3u_api_stream( stream_req: ApiStreamRequest<'_>, // _addr: &std::net::SocketAddr, ) -> impl axum::response::IntoResponse + Send { - let (user, target) = try_option_bad_request!(get_user_target_by_credentials(stream_req.username, stream_req.password, api_req, app_state), false, format!("Could not find any user {}", stream_req.username)); + let (user, target) = try_option_bad_request!(get_user_target_by_credentials(stream_req.username, stream_req.password, api_req, app_state), false, format!("Could not find any user for m3u stream {}", stream_req.username)); if user.permission_denied(app_state) { return create_custom_video_stream_response(&app_state.app_config, CustomVideoStreamType::UserAccountExpired).into_response(); } diff --git a/backend/src/api/endpoints/xtream_api.rs b/backend/src/api/endpoints/xtream_api.rs index fc7aed630..08ec0726f 100644 --- a/backend/src/api/endpoints/xtream_api.rs +++ b/backend/src/api/endpoints/xtream_api.rs @@ -178,7 +178,7 @@ async fn xtream_player_api_stream( api_req: &UserApiRequest, stream_req: ApiStreamRequest<'_>, ) -> impl IntoResponse + Send { - let (user, target) = try_option_bad_request!(get_user_target_by_credentials(stream_req.username, stream_req.password, api_req, app_state), false, format!("Could not find any user {}", stream_req.username)); + let (user, target) = try_option_bad_request!(get_user_target_by_credentials(stream_req.username, stream_req.password, api_req, app_state), false, format!("Could not find any user for xc stream {}", stream_req.username)); if user.permission_denied(app_state) { return create_custom_video_stream_response(&app_state.app_config, CustomVideoStreamType::UserAccountExpired).into_response(); } @@ -464,7 +464,7 @@ async fn xtream_player_api_resource( app_state: &Arc, resource_req: ApiStreamRequest<'_>, ) -> impl IntoResponse { - let (user, target) = try_option_bad_request!(get_user_target_by_credentials(resource_req.username, resource_req.password, api_req, app_state), false, format!("Could not find any user {}", resource_req.username)); + let (user, target) = try_option_bad_request!(get_user_target_by_credentials(resource_req.username, resource_req.password, api_req, app_state), false, format!("Could not find any user xc resource {}", resource_req.username)); if user.permission_denied(app_state) { return axum::http::StatusCode::FORBIDDEN.into_response(); } diff --git a/backend/src/api/main_api.rs b/backend/src/api/main_api.rs index 242b5764f..60399f493 100644 --- a/backend/src/api/main_api.rs +++ b/backend/src/api/main_api.rs @@ -7,11 +7,11 @@ use crate::api::endpoints::xmltv_api::xmltv_api_register; use crate::api::endpoints::xtream_api::xtream_api_register; use crate::api::model::active_provider_manager::ActiveProviderManager; use crate::api::model::active_user_manager::ActiveUserManager; -use crate::api::model::app_state::{create_cache, create_http_client, AppState, HdHomerunAppState}; +use crate::api::model::app_state::{create_cache, create_http_client, AppState, CancelTokens, HdHomerunAppState}; use crate::api::model::download::DownloadQueue; use crate::api::model::streams::shared_stream_manager::SharedStreamManager; -use crate::api::scheduler::start_scheduler; -use crate::model::{AppConfig, Config, ProcessTargets, RateLimitConfig, ScheduleConfig}; +use crate::api::scheduler::{exec_scheduler}; +use crate::model::{AppConfig, Config, ProcessTargets, RateLimitConfig}; use crate::model::{Healthcheck}; use crate::processing::processor::playlist; use log::{error, info}; @@ -20,6 +20,7 @@ use std::path::PathBuf; use std::sync::Arc; use arc_swap::{ArcSwap, ArcSwapOption}; use axum::Router; +use tokio_util::sync::CancellationToken; use tower_governor::key_extractor::SmartIpKeyExtractor; use crate::api::api_utils::{get_build_time, get_server_time}; use crate::api::config_watch::exec_config_watch; @@ -50,7 +51,7 @@ async fn healthcheck() -> impl axum::response::IntoResponse { axum::Json(create_healthcheck()) } -fn create_shared_data(app_config: &Arc) -> AppState { +fn create_shared_data(app_config: &Arc, forced_targets: &Arc) -> AppState { let config = app_config.config.load(); let cache = create_cache(&config); let active_users = Arc::new(ActiveUserManager::new(&config)); @@ -58,6 +59,7 @@ fn create_shared_data(app_config: &Arc) -> AppState { let client = create_http_client(app_config); AppState { + forced_targets: Arc::new(ArcSwap::new(Arc::clone(forced_targets))), app_config: Arc::clone(app_config), http_client: Arc::new(ArcSwap::from_pointee(client)), downloads: Arc::new(DownloadQueue::new()), @@ -65,6 +67,7 @@ fn create_shared_data(app_config: &Arc) -> AppState { shared_stream_manager: Arc::new(SharedStreamManager::new()), active_users, active_provider, + cancel_tokens: Arc::new(ArcSwap::from_pointee(CancelTokens::default())), } } @@ -80,50 +83,6 @@ fn exec_update_on_boot(client: Arc, cfg: &Arc, targe } -fn get_process_targets(cfg: &Arc, process_targets: &Arc, exec_targets: Option<&Vec>) -> Arc { - let sources = cfg.sources.load(); - if let Ok(user_targets) = sources.validate_targets(exec_targets) { - if user_targets.enabled { - if !process_targets.enabled { - return Arc::new(user_targets); - } - - let inputs: Vec = user_targets.inputs.iter() - .filter(|&id| process_targets.inputs.contains(id)) - .copied() - .collect(); - let targets: Vec = user_targets.targets.iter() - .filter(|&id| process_targets.inputs.contains(id)) - .copied() - .collect(); - return Arc::new(ProcessTargets { - enabled: user_targets.enabled, - inputs, - targets, - }); - } - } - Arc::clone(process_targets) -} - -fn exec_scheduler(client: &Arc, cfg: &Arc, targets: &Arc) { - let config = cfg.config.load(); - let schedules: Vec = if let Some(schedules) = &config.schedules { - schedules.clone() - } else { - vec![] - }; - for schedule in schedules { - let expression = schedule.schedule.to_string(); - let exec_targets = get_process_targets(cfg, targets, schedule.targets.as_ref()); - let cfg_clone = Arc::clone(cfg); - let http_client = Arc::clone(client); - tokio::spawn(async move { - start_scheduler(http_client, expression.as_str(), cfg_clone, exec_targets).await; - }); - } -} - fn is_web_auth_enabled(cfg: &Arc, web_ui_enabled: bool) -> bool { if web_ui_enabled { if let Some(web_auth) = &cfg.web_ui.as_ref().and_then(|c| c.auth.as_ref()) { @@ -149,7 +108,8 @@ fn create_compression_layer() -> tower_http::compression::CompressionLayer { .zstd(true) } -fn start_hdhomerun(app_config: &Arc, app_state: &Arc, infos: &mut Vec) { +pub(in crate::api) fn start_hdhomerun(app_config: &Arc, app_state: &Arc, infos: &mut Vec, + cancel_token: &CancellationToken) { let config = app_config.config.load(); let host = config.api.host.to_string(); let guard = app_config.hdhomerun.load(); @@ -163,6 +123,7 @@ fn start_hdhomerun(app_config: &Arc, app_state: &Arc, infos let device_clone = Arc::new(device.clone()); let basic_auth = hdhomerun.auth; infos.push(format!("HdHomeRun Server '{}' running: http://{host}:{port}", device.name)); + let c_token = cancel_token.clone(); tokio::spawn(async move { let router = axum::Router::>::new() .layer(create_cors_layer()) @@ -177,7 +138,7 @@ fn start_hdhomerun(app_config: &Arc, app_state: &Arc, infos match tokio::net::TcpListener::bind(format!("{}:{}", app_host.clone(), port)).await { Ok(listener) => { - serve(listener, router).await; + serve(listener, router, Some(c_token)).await; // if let Err(err) = axum::serve(listener, router.into_make_service_with_connect_info::()).into_future().await { // error!("{err}"); // } @@ -209,11 +170,16 @@ pub async fn start_server(app_config: Arc, targets: Arc, targets: Arc, targets: Arc = router.with_state(shared_data.clone()); let listener = tokio::net::TcpListener::bind(format!("{host}:{port}")).await?; - serve(listener, router).await + serve(listener, router, None).await; + Ok(()) //axum::serve(listener, router.into_make_service_with_connect_info::()).into_future().await } diff --git a/backend/src/api/model/app_state.rs b/backend/src/api/model/app_state.rs index 8ed35838e..e97fa4423 100644 --- a/backend/src/api/model/app_state.rs +++ b/backend/src/api/model/app_state.rs @@ -1,22 +1,48 @@ -use tokio::sync::{Mutex}; -use std::sync::Arc; -use std::time::Duration; +use crate::api::model::active_provider_manager::ActiveProviderManager; +use crate::api::model::active_user_manager::ActiveUserManager; +use crate::api::model::download::DownloadQueue; +use crate::api::model::streams::shared_stream_manager::SharedStreamManager; +use crate::api::scheduler::exec_scheduler; +use crate::model::{AppConfig, Config, HdHomeRunConfig, HdHomeRunDeviceConfig, ProcessTargets, ScheduleConfig}; +use crate::tools::lru_cache::LRUResourceCache; +use crate::utils::request::create_client; use arc_swap::{ArcSwap, ArcSwapAny}; use log::error; use reqwest::Client; use shared::error::TuliproxError; use shared::model::UserConnectionPermission; -use crate::api::model::active_provider_manager::ActiveProviderManager; -use crate::api::model::active_user_manager::ActiveUserManager; -use crate::api::model::download::DownloadQueue; -use crate::api::model::streams::shared_stream_manager::SharedStreamManager; -use crate::model::{AppConfig, Config, HdHomeRunDeviceConfig}; -use crate::tools::lru_cache::LRUResourceCache; -use crate::utils::request::create_client; +use shared::utils::small_vecs_equal_unordered; +use std::sync::Arc; +use std::time::Duration; +use tokio::sync::Mutex; +use tokio_util::sync::CancellationToken; + +pub(in crate::api) struct UpdateChanges { + scheduler: bool, + hdhomerun: bool, +} + +pub async fn update_app_state(app_state: &Arc, config: Config) -> Result<(), TuliproxError> { + let updates = app_state.set_config(config).await?; + start_services(app_state, &updates); + Ok(()) +} + +fn start_services(app_state: &Arc, changes: &UpdateChanges) { + if changes.scheduler { + exec_scheduler(&Arc::clone(&app_state.http_client.load()), &app_state.app_config, + &app_state.forced_targets.load(), &app_state.cancel_tokens.load().scheduler); + } + + if changes.hdhomerun && app_state.app_config.api_proxy.load().is_some() { + let mut infos = Vec::new(); + crate::api::main_api::start_hdhomerun(&app_state.app_config, app_state, &mut infos, &app_state.cancel_tokens.load().hdhomerun); + } +} pub fn create_http_client(app_config: &AppConfig) -> Client { let mut builder = create_client(app_config).http1_only(); - let config = app_config.config.load();// because of RAII connection dropping + let config = app_config.config.load(); // because of RAII connection dropping if config.connect_timeout_secs > 0 { builder = builder.connect_timeout(Duration::from_secs(u64::from(config.connect_timeout_secs))); } @@ -44,27 +70,52 @@ pub fn create_cache(config: &Config) -> Option>> { None } +pub struct CancelTokens { + pub(crate) scheduler: CancellationToken, + pub(crate) hdhomerun: CancellationToken, +} +impl Default for CancelTokens { + fn default() -> Self { + Self { + scheduler: CancellationToken::new(), + hdhomerun: CancellationToken::new(), + } + } +} + +macro_rules! change_detect { + ($fn_name:ident, $a:expr, $b: expr) => { + match ($a, $b) { + (None, None) => false, + (Some(_), None) | + (None, Some(_)) => true, + (Some(o), Some(n)) => $fn_name(o, n), + } + }; +} #[derive(Clone)] pub struct AppState { + pub forced_targets: Arc>, // as program arguments pub app_config: Arc, - pub http_client: Arc>, + pub http_client: Arc>, pub downloads: Arc, pub cache: Arc>>>>, pub shared_stream_manager: Arc, pub active_users: Arc, pub active_provider: Arc, + pub cancel_tokens: Arc>, } impl AppState { - - pub async fn set_config(&self, config: Config) -> Result<(), TuliproxError> { + pub(in crate::api::model) async fn set_config(&self, config: Config) -> Result { + let changes = self.detect_changes(&config); config.update_runtime(); self.active_users.update_config(&config); self.active_provider.update_config(&self.app_config).await; self.app_config.set_config(config)?; self.update_config().await; - Ok(()) + Ok(changes) } async fn update_config(&self) { @@ -100,6 +151,65 @@ impl AppState { pub async fn get_connection_permission(&self, username: &str, max_connections: u32) -> UserConnectionPermission { self.active_users.connection_permission(username, max_connections).await } + + fn detect_changes(&self, config: &Config) -> UpdateChanges { + let old_config = self.app_config.config.load(); + let changed_schedules = change_detect!(schedules_changed, old_config.schedules.as_ref(), config.schedules.as_ref()); + let changed_hdhomerun = change_detect!(hdhomerun_changed, old_config.hdhomerun.as_ref(), config.hdhomerun.as_ref()); + + if changed_schedules || changed_hdhomerun { + let cancel_tokens = self.cancel_tokens.load(); + if changed_schedules { + cancel_tokens.scheduler.cancel(); + } + if changed_hdhomerun { + cancel_tokens.hdhomerun.cancel(); + } + + let tokens = CancelTokens { + scheduler: if changed_schedules { CancellationToken::default() } else { cancel_tokens.scheduler.clone() }, + hdhomerun: if changed_hdhomerun { CancellationToken::default() } else { cancel_tokens.hdhomerun.clone() }, + }; + self.cancel_tokens.store(Arc::new(tokens)); + } + UpdateChanges { + scheduler: changed_schedules, + hdhomerun: changed_hdhomerun, + } + } +} + +fn schedules_changed(a: &[ScheduleConfig], b: &[ScheduleConfig]) -> bool { + if a.len() != b.len() { + return true; + } + for schedule in a { + if let Some(found) = b.iter().find(|&s| s.schedule == schedule.schedule) { + match (schedule.targets.as_ref(), found.targets.as_ref()) { + (None, None) => return false, + (Some(_targets), None) | + (None, Some(_targets)) => return true, + (Some(a_targets), Some(b_targets)) => { + if !small_vecs_equal_unordered(a_targets, b_targets) { + return true; + } + } + } + } else { + return true; + } + } + false +} + +fn hdhomerun_changed(a: &HdHomeRunConfig, b: &HdHomeRunConfig) -> bool { + if a.enabled != b.enabled || a.auth != b.auth { + return true; + } + if !small_vecs_equal_unordered(a.devices.as_ref(), b.devices.as_ref()) { + return true; + } + false } #[derive(Clone)] diff --git a/backend/src/api/scheduler.rs b/backend/src/api/scheduler.rs index c9f7b275a..bea87ebd4 100644 --- a/backend/src/api/scheduler.rs +++ b/backend/src/api/scheduler.rs @@ -5,7 +5,8 @@ use chrono::{DateTime, FixedOffset, Local}; use cron::Schedule; use crate::utils::{exit}; use log::{error}; -use crate::model::{AppConfig, ProcessTargets}; +use tokio_util::sync::CancellationToken; +use crate::model::{AppConfig, ProcessTargets, ScheduleConfig}; use crate::processing::processor::playlist::exec_processing; pub fn datetime_to_instant(datetime: DateTime) -> Instant { @@ -24,22 +25,81 @@ pub fn datetime_to_instant(datetime: DateTime) -> Instant { Instant::now() + duration_until } -pub async fn start_scheduler(client: Arc, expression: &str, config: Arc, targets: Arc) -> ! { +pub fn exec_scheduler(client: &Arc, cfg: &Arc, targets: &Arc, + cancel: &CancellationToken) { + let config = cfg.config.load(); + let schedules: Vec = if let Some(schedules) = &config.schedules { + schedules.clone() + } else { + vec![] + }; + for schedule in schedules { + let expression = schedule.schedule.to_string(); + let exec_targets = get_process_targets(cfg, targets, schedule.targets.as_ref()); + let cfg_clone = Arc::clone(cfg); + let http_client = Arc::clone(client); + let cancel_token = cancel.clone(); + tokio::spawn(async move { + start_scheduler(http_client, expression.as_str(), cfg_clone, exec_targets, cancel_token).await; + }); + } +} + +async fn start_scheduler(client: Arc, expression: &str, config: Arc, + targets: Arc, cancel: CancellationToken) { match Schedule::from_str(expression) { Ok(schedule) => { let offset = *Local::now().offset(); loop { let mut upcoming = schedule.upcoming(offset).take(1); if let Some(datetime) = upcoming.next() { - tokio::time::sleep_until(tokio::time::Instant::from(datetime_to_instant(datetime))).await; - exec_processing(Arc::clone(&client), Arc::clone(&config), Arc::clone(&targets)).await; - } + tokio::select! { + () = tokio::time::sleep_until(tokio::time::Instant::from(datetime_to_instant(datetime))) => { + exec_processing(Arc::clone(&client), Arc::clone(&config), Arc::clone(&targets)).await; + } + () = cancel.cancelled() => { + break; + } + } + } } } Err(err) => exit!("Failed to start scheduler: {}", err) } } + +fn get_process_targets(cfg: &Arc, process_targets: &Arc, exec_targets: Option<&Vec>) -> Arc { + let sources = cfg.sources.load(); + if let Ok(user_targets) = sources.validate_targets(exec_targets) { + if user_targets.enabled { + if !process_targets.enabled { + return Arc::new(user_targets); + } + + let inputs: Vec = user_targets.inputs.iter() + .filter(|&id| process_targets.inputs.contains(id)) + .copied() + .collect(); + let targets: Vec = user_targets.targets.iter() + .filter(|&id| process_targets.inputs.contains(id)) + .copied() + .collect(); + let target_names: Vec = user_targets.target_names.iter() + .filter(|&name| process_targets.target_names.contains(name)) + .cloned() + .collect(); + return Arc::new(ProcessTargets { + enabled: user_targets.enabled, + inputs, + targets, + target_names + }); + } + } + Arc::clone(process_targets) +} + #[cfg(test)] mod tests { use std::str::FromStr; diff --git a/backend/src/api/serve.rs b/backend/src/api/serve.rs index 14dde54b0..f3ebe1a5c 100644 --- a/backend/src/api/serve.rs +++ b/backend/src/api/serve.rs @@ -14,6 +14,7 @@ use std::net::SocketAddr; use std::pin::pin; use std::time::Duration; use tokio::sync::watch; +use tokio_util::sync::CancellationToken; use tower::{Service, ServiceExt}; #[derive(Debug)] @@ -35,48 +36,41 @@ impl axum::extract::connect_info::Connected for SocketAddr { } } -pub async fn serve(listener: tokio::net::TcpListener, router: axum::Router<()>) -> ! { +pub async fn serve(listener: tokio::net::TcpListener, router: axum::Router<()>, + cancel_token: Option) { let (signal_tx, _signal_rx) = watch::channel(()); let (_close_tx, close_rx) = watch::channel(()); let mut make_service = router.into_make_service_with_connect_info::(); - loop { - let Ok((socket, remote_addr)) = listener.accept().await else { continue }; - - let Ok(tcp_stream_std) = socket.into_std() else { continue; }; - tcp_stream_std.set_nonblocking(true).ok(); // this is not necessary - - // Configure keep alive with socket2 - let sock_ref = SockRef::from(&tcp_stream_std); - - let keep_alive_first_probe = 10; - let keep_alive_interval = 5; - - let mut keepalive = TcpKeepalive::new(); - keepalive = keepalive.with_time(Duration::from_secs(keep_alive_first_probe)) // Time until the first keepalive probe (idle time) - .with_interval(Duration::from_secs(keep_alive_interval)); // Interval between keep alives - #[cfg(not(target_os = "windows"))] - { - let keep_alive_retries = 3; - keepalive = keepalive.with_retries(keep_alive_retries); // Number of failed probes before the connection is closed + match cancel_token { + Some(token) => { + loop { + tokio::select! { + () = token.cancelled() => { + break; + } + accept_result = listener.accept() => { + let Ok((socket, remote_addr)) = accept_result else { continue }; + handle_connection(&mut make_service, &signal_tx, &close_rx, socket, remote_addr).await; + } + } + } } - - if let Err(e) = sock_ref.set_tcp_keepalive(&keepalive) { - error!("Failed to set keepalive for {remote_addr}: {e}"); + None => { + loop { + let Ok((socket, remote_addr)) = listener.accept().await else { continue }; + handle_connection(&mut make_service, &signal_tx, &close_rx, socket, remote_addr).await; + } } - - let Ok(socket) = tokio::net::TcpStream::from_std(tcp_stream_std) else { continue; }; - - let io = TokioIo::new(socket); - handle_connection(&mut make_service, &signal_tx, &close_rx, io, remote_addr).await; } } + async fn handle_connection( make_service: &mut M, signal_tx: &watch::Sender<()>, close_rx: &watch::Receiver<()>, - io: TokioIo, + socket: tokio::net::TcpStream, remote_addr: SocketAddr, ) where @@ -85,6 +79,31 @@ where S: Service + Clone + Send + 'static, S::Future: Send, { + let Ok(tcp_stream_std) = socket.into_std() else { return; }; + tcp_stream_std.set_nonblocking(true).ok(); // this is not necessary + + // Configure keep alive with socket2 + let sock_ref = SockRef::from(&tcp_stream_std); + + let keep_alive_first_probe = 10; + let keep_alive_interval = 5; + + let mut keepalive = TcpKeepalive::new(); + keepalive = keepalive.with_time(Duration::from_secs(keep_alive_first_probe)) // Time until the first keepalive probe (idle time) + .with_interval(Duration::from_secs(keep_alive_interval)); // Interval between keep alives + #[cfg(not(target_os = "windows"))] + { + let keep_alive_retries = 3; + keepalive = keepalive.with_retries(keep_alive_retries); // Number of failed probes before the connection is closed + } + + if let Err(e) = sock_ref.set_tcp_keepalive(&keepalive) { + error!("Failed to set keepalive for {remote_addr}: {e}"); + } + + let Ok(socket) = tokio::net::TcpStream::from_std(tcp_stream_std) else { return; }; + + let io = TokioIo::new(socket); trace!("connection {remote_addr:?} accepted"); make_service diff --git a/backend/src/main.rs b/backend/src/main.rs index 393df7002..bb23522e3 100644 --- a/backend/src/main.rs +++ b/backend/src/main.rs @@ -10,7 +10,7 @@ mod modules; include_modules!(); use crate::auth::generate_password; -use crate::model::{AppConfig, Config, Healthcheck, HealthcheckConfig, ProcessTargets}; +use crate::model::{AppConfig, Config, Healthcheck, HealthcheckConfig, ProcessTargets, SourcesConfig}; use crate::processing::processor::playlist; use crate::utils::{config_file_reader, resolve_env_var}; use crate::utils::request::{create_client}; @@ -19,6 +19,8 @@ use clap::{Parser}; use log::{error, info}; use std::fs::File; use std::sync::Arc; +use arc_swap::access::Access; +use arc_swap::ArcSwap; use shared::model::ConfigPaths; use crate::utils::init_logger; @@ -106,8 +108,9 @@ fn main() { } let app_config = utils::read_initial_app_config(&mut config_paths, true, true, args.server).unwrap_or_else(|err| exit!("{}", err)); + print_info(&app_config); - let sources = app_config.sources.load(); + let sources = > as Access>::load(&app_config.sources); let targets = sources.validate_targets(args.target.as_ref()).unwrap_or_else(|err| exit!("{}", err)); let rt = tokio::runtime::Runtime::new().unwrap(); @@ -120,6 +123,28 @@ fn main() { }); } +fn print_info(app_config: &AppConfig) { + let config = > as Access>::load(&app_config.config); + let paths = > as Access>::load(&app_config.paths); + info!("Current time: {}", chrono::offset::Local::now().format("%Y-%m-%d %H:%M:%S")); + info!("Temp dir: {}", tempfile::env::temp_dir().display()); + info!("Working dir: {:?}", &config.working_dir); + info!("Config dir: {:?}", &paths.config_path); + info!("Config file: {:?}", &paths.config_file_path); + info!("Source file: {:?}", &paths.sources_file_path); + info!("Api Proxy File: {:?}", &paths.api_proxy_file_path); + info!("Mapping file: {:?}", &paths.mapping_file_path.as_ref().map_or_else(|| "not used", |v| v.as_str())); + + if let Some(cache) = config.reverse_proxy.as_ref().and_then(|r| r.cache.as_ref()) { + if cache.enabled { + info!("Cache dir: {}", cache.dir); + } + } + if let Some(resource_path) = paths.custom_stream_response_path.as_ref() { + info!("Resource path: {resource_path}"); + } +} + fn get_file_paths(args: &Args) -> ConfigPaths { let config_path: String = utils::resolve_directory_path(&resolve_env_var(&args.config_path.as_ref().map_or_else(utils::get_default_config_path, ToString::to_string))); let config_file: String = resolve_env_var(&args.config_file.as_ref().map_or_else(|| utils::get_default_config_file_path(&config_path), ToString::to_string)); diff --git a/backend/src/model/config/api_proxy.rs b/backend/src/model/config/api_proxy.rs index ae6f986f5..c923671f9 100644 --- a/backend/src/model/config/api_proxy.rs +++ b/backend/src/model/config/api_proxy.rs @@ -188,7 +188,7 @@ impl ApiProxyConfig { .find(|credential| credential.username == username) .cloned(); if result.is_none() && (username != "test" || username != "api") { - debug!("Could not find any user {username}"); + debug!("Could not find any user credentials for: {username}"); } result } diff --git a/backend/src/model/config/app.rs b/backend/src/model/config/app.rs index aa659d71f..b3c3b85c0 100644 --- a/backend/src/model/config/app.rs +++ b/backend/src/model/config/app.rs @@ -6,7 +6,7 @@ use std::path::{Path, PathBuf}; use std::sync::Arc; use arc_swap::{ArcSwap, ArcSwapOption}; use arc_swap::access::Access; -use log::{debug, error, info}; +use log::{debug, error}; use rand::Rng; use shared::create_tuliprox_error_result; use shared::error::{TuliproxError, TuliproxErrorKind}; @@ -46,7 +46,6 @@ impl AppConfig { pub fn set_config(&self, config: Config) -> Result<(), TuliproxError> { self.config.store(Arc::new(config)); self.prepare_custom_stream_response(); - self.print_info(); Ok(()) } @@ -377,27 +376,6 @@ impl AppConfig { self.get_server_info(server_info_name) } - pub fn print_info(&self) { - let config = > as Access>::load(&self.config); - let paths = > as Access>::load(&self.paths); - info!("Current time: {}", chrono::offset::Local::now().format("%Y-%m-%d %H:%M:%S")); - info!("Temp dir: {}", tempfile::env::temp_dir().display()); - info!("Working dir: {:?}", &config.working_dir); - info!("Config dir: {:?}", &paths.config_path); - info!("Config file: {:?}", &paths.config_file_path); - info!("Source file: {:?}", &paths.sources_file_path); - info!("Api Proxy File: {:?}", &paths.api_proxy_file_path); - info!("Mapping file: {:?}", &paths.mapping_file_path.as_ref().map_or_else(|| "not used", |v| v.as_str())); - - if let Some(cache) = config.reverse_proxy.as_ref().and_then(|r| r.cache.as_ref()) { - if cache.enabled { - info!("Cache dir: {}", cache.dir); - } - } - if let Some(resource_path) = paths.custom_stream_response_path.as_ref() { - info!("Resource path: {resource_path}"); - } - } } diff --git a/backend/src/model/config/hdhomerun.rs b/backend/src/model/config/hdhomerun.rs index b38e38ce1..d23a50e67 100644 --- a/backend/src/model/config/hdhomerun.rs +++ b/backend/src/model/config/hdhomerun.rs @@ -1,7 +1,7 @@ use shared::model::{HdHomeRunConfigDto, HdHomeRunDeviceConfigDto}; use crate::model::macros; -#[derive(Debug, Clone)] +#[derive(Debug, Clone, PartialEq, Eq)] pub struct HdHomeRunDeviceConfig { pub friendly_name: String, pub manufacturer: String, diff --git a/backend/src/model/config/source.rs b/backend/src/model/config/source.rs index 1d2e905c5..f2c410bf0 100644 --- a/backend/src/model/config/source.rs +++ b/backend/src/model/config/source.rs @@ -78,6 +78,7 @@ impl SourcesConfig { let mut enabled = true; let mut inputs: Vec = vec![]; let mut targets: Vec = vec![]; + let mut target_names: Vec = vec![]; if let Some(user_targets) = target_args { let mut check_targets: HashMap = user_targets.iter().map(|t| (t.to_lowercase(), 0)).collect(); for source in &self.sources { @@ -87,6 +88,7 @@ impl SourcesConfig { let key = user_target.to_lowercase(); if target.name.eq_ignore_ascii_case(key.as_str()) { targets.push(target.id); + target_names.push(target.name.to_string()); target_added = true; if let Some(value) = check_targets.get(key.as_str()) { check_targets.insert(key, value + 1); @@ -113,6 +115,7 @@ impl SourcesConfig { enabled, inputs, targets, + target_names, }) } diff --git a/backend/src/model/config/target.rs b/backend/src/model/config/target.rs index 63c152d75..e277c5390 100644 --- a/backend/src/model/config/target.rs +++ b/backend/src/model/config/target.rs @@ -15,6 +15,7 @@ pub struct ProcessTargets { pub enabled: bool, pub inputs: Vec, pub targets: Vec, + pub target_names: Vec, } impl ProcessTargets { diff --git a/backend/src/utils/file/config_reader.rs b/backend/src/utils/file/config_reader.rs index 7f3baab02..6da877dfb 100644 --- a/backend/src/utils/file/config_reader.rs +++ b/backend/src/utils/file/config_reader.rs @@ -1,6 +1,6 @@ use crate::model::Config; use crate::model::{ApiProxyConfig, AppConfig, SourcesConfig}; -use crate::utils; +use crate::{print_info, utils}; use crate::utils::file_reader; use crate::utils::sys_utils::exit; use crate::utils::{open_file, read_mappings_file, EnvResolvingReader, FileLockManager}; @@ -160,7 +160,7 @@ pub fn read_initial_app_config(paths: &mut ConfigPaths, encrypt_secret: Default::default(), }; app_config.prepare(include_computed)?; - app_config.print_info(); + print_info(&app_config); if let Some(mappings_file) = &paths.mapping_file_path { match utils::read_mappings(mappings_file.as_str(), resolve_env) { diff --git a/shared/src/utils/string_utils.rs b/shared/src/utils/string_utils.rs index 0adc27586..075d97183 100644 --- a/shared/src/utils/string_utils.rs +++ b/shared/src/utils/string_utils.rs @@ -42,6 +42,20 @@ pub fn generate_random_string(length: usize) -> String { random_string } +// compare 2 small vecs without HashSet +pub fn small_vecs_equal_unordered(a: &[T], b: &[T]) -> bool { + if a.len() != b.len() { + return false; + } + + for item in a { + if !b.iter().any(|x| x == item) { + return false; + } + } + true +} + pub fn get_non_empty_str<'a>(first: &'a str, second: &'a str, third: &'a str) -> &'a str { if !first.is_empty() { first