diff --git a/CHANGELOG.md b/CHANGELOG.md index c595ee8b0..1e62c193f 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,5 @@ # Changelog -# 2.2.6 (2025-03-xx) +# 2.2.6 (2025-04-xx) - !BREAKING CHANGE! bandwidth `throttle_kbps` attribute for `reverse_proxy.stream` in `config.yml` is now `throttle` and supports units. Allowed units are `KB/s`,`MB/s`,`KiB/s`,`MiB/s`,`kbps`,`mbps`,`Mibps`. Default unit is `kbps`. @@ -28,6 +28,7 @@ web_ui: userfile: user.txt ``` - user has now the attribute `ui_enabled` to disable/enable web_ui for user. +- epg processing optimization, auto guessing/assigning epg id's # 2.2.5 (2025-03-27) - fixed web ui playlist regexp search diff --git a/src/api/api_utils.rs b/src/api/api_utils.rs index 80024e257..5ac94c915 100644 --- a/src/api/api_utils.rs +++ b/src/api/api_utils.rs @@ -76,6 +76,7 @@ use crate::api::model::streams::throttled_stream::ThrottledStream; use crate::tools::atomic_once_flag::AtomicOnceFlag; use crate::utils::default_utils::default_grace_period_millis; +#[allow(clippy::missing_panics_doc)] pub async fn serve_file(file_path: &Path, mime_type: mime::Mime) -> impl axum::response::IntoResponse + Send { if file_path.exists() { return match tokio::fs::File::open(file_path).await { @@ -165,8 +166,7 @@ fn get_stream_alternative_url(stream_url: &str, input: &ConfigInput, alias_input let modified = stream_url.replace(&input_user_info.base_url, &alt_input_user_info.base_url); let modified = modified.replace(&input_user_info.username, &alt_input_user_info.username); - let modified = modified.replace(&input_user_info.password, &alt_input_user_info.password); - modified + modified.replace(&input_user_info.password, &alt_input_user_info.password) } type StreamUrl = String; @@ -221,12 +221,13 @@ fn get_streaming_options(app_state: &AppState, stream_url: &str, input_opt: Opti } ProviderAllocation::Available(provider) | ProviderAllocation::GracePeriod(provider) => { - let (provider, url) = if provider.id != input.id { - (provider.name.to_string(), get_stream_alternative_url(stream_url, input, &provider)) - } else { + let (provider, url) = if provider.id == input.id { (input.name.to_string(), stream_url.to_string()) + } else { + (provider.name.to_string(), get_stream_alternative_url(stream_url, input, provider)) }; + if matches!(allocation, ProviderAllocation::Available(_)) { StreamingOption::AvailableStream(Some(provider), url) } else { @@ -244,7 +245,7 @@ async fn create_stream_response_details(app_state: &AppState, stream_options: &S req_headers: &HeaderMap, input_opt: Option<&ConfigInput>, item_type: PlaylistItemType, share_stream: bool) -> StreamDetails { let (stream_response_params, input_headers) = get_streaming_options(app_state, stream_url, input_opt); - let config_grace_period_millis = app_state.config.reverse_proxy.as_ref().and_then(|r| r.stream.as_ref()).map(|s| s.grace_period_millis).unwrap_or_else(default_grace_period_millis); + let config_grace_period_millis = app_state.config.reverse_proxy.as_ref().and_then(|r| r.stream.as_ref()).map_or_else(default_grace_period_millis, |s| s.grace_period_millis); let grace_period_millis = if config_grace_period_millis > 0 && matches!(stream_response_params, StreamingOption::GracePeriodStream(_, _)) { config_grace_period_millis } else { 0 }; match stream_response_params { StreamingOption::CustomStream(provider_stream) => { @@ -262,11 +263,11 @@ async fn create_stream_response_details(app_state: &AppState, stream_options: &S let parsed_url = Url::parse(&request_url); let ((stream, stream_info), reconnect_flag) = if let Ok(url) = parsed_url { if stream_options.pipe_provider_stream { - (provider_stream::get_provider_pipe_stream(app_state, &url, req_headers, input_headers, item_type).await, None) + (provider_stream::get_provider_pipe_stream(app_state, &url, req_headers, input_headers.as_ref(), item_type).await, None) } else { - let buffer_stream_options = BufferStreamOptions::new(item_type, share_stream, &stream_options); + let buffer_stream_options = BufferStreamOptions::new(item_type, share_stream, stream_options); let reconnect_flag = buffer_stream_options.get_reconnect_flag_clone(); - (provider_stream::get_provider_reconnect_buffered_stream(app_state, &url, req_headers, input_headers, buffer_stream_options).await, + (provider_stream::get_provider_reconnect_buffered_stream(app_state, &url, req_headers, input_headers.as_ref(), buffer_stream_options).await, Some(reconnect_flag)) } } else { @@ -320,12 +321,12 @@ pub async fn stream_response(app_state: &AppState, let stream_options = get_stream_options(app_state); let stream_details = - create_stream_response_details(app_state, &stream_options, &stream_url, req_headers, input, item_type, share_stream).await; + create_stream_response_details(app_state, &stream_options, stream_url, req_headers, input, item_type, share_stream).await; if stream_details.has_stream() { // let content_length = get_stream_content_length(provider_response.as_ref()); - let provider_response = stream_details.stream_info.as_ref().map_or(None, |(h, sc)| Some((h.clone(), sc.clone()))); - let stream = ActiveClientStream::new(stream_details, app_state, &user).await; + let provider_response = stream_details.stream_info.as_ref().map(|(h, sc)| (h.clone(), *sc)); + let stream = ActiveClientStream::new(stream_details, app_state, user).await; let stream_resp = if share_stream { // Shared Stream response let shared_headers = provider_response.as_ref().map_or_else(Vec::new, |(h, _)| h.clone()); @@ -381,7 +382,7 @@ async fn shared_stream_response(app_state: &AppState, stream_url: &str, user: &P if let Some(headers) = app_state.shared_stream_manager.get_shared_state_headers(stream_url).await { let (status_code, header_map) = get_stream_response_with_headers(Some((headers.clone(), StatusCode::OK))); let stream_details = StreamDetails::from_stream(stream); - let stream = ActiveClientStream::new(stream_details, app_state, &user).await.boxed(); + let stream = ActiveClientStream::new(stream_details, app_state, user).await.boxed(); let mut response = axum::response::Response::builder() .status(status_code); for (key, value) in &header_map { diff --git a/src/api/endpoints/download_api.rs b/src/api/endpoints/download_api.rs index e31e8c8ec..54c04e850 100644 --- a/src/api/endpoints/download_api.rs +++ b/src/api/endpoints/download_api.rs @@ -21,7 +21,7 @@ async fn download_file(active: Arc>>, client: &reqwe match fs::create_dir_all(&file_download.file_dir) { Ok(()) => { if let Some(file_path_str) = file_download.file_path.to_str() { - info!("Downloading {}", file_path_str); + info!("Downloading {file_path_str}"); match File::create(&file_download.file_path) { Ok(mut file) => { let mut downloaded: u64 = 0; @@ -39,7 +39,7 @@ async fn download_file(active: Arc>>, client: &reqwe } } else { let megabytes = request::bytes_to_megabytes(downloaded); - info!("Downloaded {}, filesize: {}MB", file_path_str, megabytes); + info!("Downloaded {file_path_str}, filesize: {megabytes}MB"); active.write().await.as_mut().unwrap().size = downloaded; return Ok(()); } diff --git a/src/api/endpoints/v1_api.rs b/src/api/endpoints/v1_api.rs index 4d9c8a639..050b32d1c 100644 --- a/src/api/endpoints/v1_api.rs +++ b/src/api/endpoints/v1_api.rs @@ -15,7 +15,7 @@ use crate::api::model::request::{PlaylistRequest, PlaylistRequestType}; use crate::auth::access_token::create_access_token; use crate::auth::authenticator::{validator_admin}; use crate::m3u_filter_error::M3uFilterError; -use crate::model::api_proxy::{ApiProxyConfig, ApiProxyServerInfo, TargetUser}; +use crate::model::api_proxy::{ApiProxyConfig, ApiProxyServerInfo, ProxyUserCredentials, TargetUser}; use crate::model::config::{validate_targets, Config, ConfigDto, ConfigInput, ConfigInputOptions, ConfigSource, ConfigTarget, InputType, TargetType}; use crate::model::playlist::{XtreamCluster, XtreamPlaylistItem}; use crate::model::xtream::PlaylistXtreamCategory; @@ -31,7 +31,7 @@ fn intern_save_config_api_proxy(backup_dir: &str, api_proxy: &ApiProxyConfig, fi match config_reader::save_api_proxy(file_path, backup_dir, api_proxy) { Ok(()) => {} Err(err) => { - error!("Failed to save api_proxy.yml {}", err.to_string()); + error!("Failed to save api_proxy.yml {err}"); return Some(err); } } @@ -42,7 +42,7 @@ fn intern_save_config_main(file_path: &str, backup_dir: &str, cfg: &ConfigDto) - match config_reader::save_main_config(file_path, backup_dir, cfg) { Ok(()) => {} Err(err) => { - error!("Failed to save config.yml {}", err.to_string()); + error!("Failed to save config.yml {err}"); return Some(err); } } @@ -77,7 +77,7 @@ async fn save_config_api_proxy_user( let mut lock = app_state.config.t_api_proxy.write().await; if let Some(api_proxy) = lock.as_mut() { api_proxy.user = users; - api_proxy.user.iter_mut().flat_map(|t| &mut t.credentials).for_each(|c| c.prepare()); + api_proxy.user.iter_mut().flat_map(|t| &mut t.credentials).for_each(ProxyUserCredentials::prepare); if api_proxy.use_user_db { if let Err(err) = store_api_user(&app_state.config, &api_proxy.user) { return (axum::http::StatusCode::INTERNAL_SERVER_ERROR, axum::Json(json!({"error": err.to_string()}))).into_response(); @@ -304,7 +304,7 @@ async fn playlist_reverse( ) -> impl axum::response::IntoResponse + Send { let access_token = create_access_token(&app_state.config.t_access_token_secret, 5); let server_name = app_state.config.web_ui.as_ref().and_then(|web_ui| web_ui.player_server.as_ref()).map_or("default", |server_name| server_name.as_str()); - let server_info = app_state.config.get_server_info(&server_name).await; + let server_info = app_state.config.get_server_info(server_name).await; let base_url = server_info.get_base_url(); format!("{base_url}/token/{access_token}/{target_id}/{}/{}", playlist_item.xtream_cluster.as_stream_type(), playlist_item.virtual_id).into_response() } @@ -315,7 +315,7 @@ async fn config( let map_input = |i: &ConfigInput| ServerInputConfig { id: i.id, name: i.name.clone(), - input_type: i.input_type.clone(), + input_type: i.input_type, url: i.url.clone(), username: i.username.clone(), password: i.password.clone(), @@ -393,7 +393,7 @@ pub fn v1_api_register(web_auth_enabled: bool, app_state: Arc, web_ui_ } let mut base_router = axum::Router::new(); - if app_state.config.web_ui.as_ref().map_or(true, |c| c.user_ui_enabled) { + if app_state.config.web_ui.as_ref().is_none_or(|c| c.user_ui_enabled) { base_router = base_router.merge(user_api_register(app_state)); } base_router.nest(&format!("{web_ui_path}/api/v1"), router) diff --git a/src/api/endpoints/web_index.rs b/src/api/endpoints/web_index.rs index 4d5a02c49..37c34dc5a 100644 --- a/src/api/endpoints/web_index.rs +++ b/src/api/endpoints/web_index.rs @@ -99,7 +99,7 @@ async fn index( let base_href = format!(r#""#); if let Some(pos) = new_content.find("") { new_content.replace_range(pos..pos + 6, &base_href); - }; + } return axum::response::Response::builder() .header("Content-Type", mime::TEXT_HTML_UTF_8.as_ref()) diff --git a/src/api/main_api.rs b/src/api/main_api.rs index eb7be0a18..77746b154 100644 --- a/src/api/main_api.rs +++ b/src/api/main_api.rs @@ -40,7 +40,7 @@ fn get_web_dir_path(web_ui_enabled: bool, web_root: &str) -> Result, targets: Arc) -> fut let mut infos = Vec::new(); let host = cfg.api.host.to_string(); let port = cfg.api.port; - let web_ui_enabled = cfg.web_ui.as_ref().map_or(false, |c| c.enabled); + let web_ui_enabled = cfg.web_ui.as_ref().is_some_and(|c| c.enabled); let web_dir_path = match get_web_dir_path(web_ui_enabled, cfg.api.web_root.as_str()) { Ok(result) => result, Err(err) => return Err(err) @@ -318,7 +318,7 @@ pub async fn start_server(cfg: Arc, targets: Arc) -> fut let mut rate_limiting = false; if let Some(rate_limiter) = app_state.config.reverse_proxy.as_ref().and_then(|r| r.rate_limit.clone()) { rate_limiting = rate_limiter.enabled; - api_router = add_rate_limiter(api_router, rate_limiter); + api_router = add_rate_limiter(api_router, &rate_limiter); } router = router.merge(api_router); @@ -342,7 +342,7 @@ pub async fn start_server(cfg: Arc, targets: Arc) -> fut } } -fn add_rate_limiter(router: Router>, rate_limit_cfg: RateLimitConfig) -> Router> { +fn add_rate_limiter(router: Router>, rate_limit_cfg: &RateLimitConfig) -> Router> { if rate_limit_cfg.enabled { let governor_conf = Arc::new(tower_governor::governor::GovernorConfigBuilder::default() .key_extractor(SmartIpKeyExtractor) diff --git a/src/api/model/active_provider_manager.rs b/src/api/model/active_provider_manager.rs index c561b3c77..9845d671d 100644 --- a/src/api/model/active_provider_manager.rs +++ b/src/api/model/active_provider_manager.rs @@ -38,7 +38,7 @@ impl ProviderConfig { url: cfg.url.clone(), username: cfg.username.clone(), password: cfg.password.clone(), - input_type: cfg.input_type.clone(), + input_type: cfg.input_type, max_connections: cfg.max_connections, priority: cfg.priority, current_connections: AtomicU16::new(0), @@ -52,7 +52,7 @@ impl ProviderConfig { url: alias.url.clone(), username: alias.username.clone(), password: alias.password.clone(), - input_type: cfg.input_type.clone(), + input_type: cfg.input_type, max_connections: alias.max_connections, priority: alias.priority, current_connections: AtomicU16::new(0), @@ -60,7 +60,7 @@ impl ProviderConfig { } pub fn get_user_info(&self) -> Option { - InputUserInfo::new(self.input_type.clone(), self.username.as_deref(), self.password.as_deref(), &self.url) + InputUserInfo::new(self.input_type, self.username.as_deref(), self.password.as_deref(), &self.url) } #[inline] @@ -259,7 +259,8 @@ impl MultiProviderLineup { ProviderPriorityGroup::MultiProviderGroup(index, pg) => { let mut idx = index.load(Ordering::SeqCst); let provider_count = pg.len(); - for _ in idx..provider_count { + let start = idx; + for _ in start..provider_count { let p = pg.get(idx).unwrap(); idx = (idx + 1) % provider_count; let result = p.try_allocate(grace); diff --git a/src/api/model/active_user_manager.rs b/src/api/model/active_user_manager.rs index 4cb8fb94d..da1d34ecf 100644 --- a/src/api/model/active_user_manager.rs +++ b/src/api/model/active_user_manager.rs @@ -31,7 +31,7 @@ impl ActiveUserManager { pub async fn connection_permission(&self, username: &str, max_connections: u32, grace_period: bool) -> UserConnectionPermission { if let Some(counter) = self.user.read().await.get(username) { let current_connections = counter.load(Ordering::SeqCst); - let extra_con = if grace_period { 1 } else { 0 }; + let extra_con = u32::from(grace_period); if current_connections < max_connections + extra_con { return UserConnectionPermission::GracePeriod; } diff --git a/src/api/model/app_state.rs b/src/api/model/app_state.rs index dba307ae4..efa5eb143 100644 --- a/src/api/model/app_state.rs +++ b/src/api/model/app_state.rs @@ -29,7 +29,7 @@ impl AppState { } pub async fn get_connection_permission(&self, username: &str, max_connections: u32) -> UserConnectionPermission { - let grace_period_millis = self.config.reverse_proxy.as_ref().and_then(|r| r.stream.as_ref()).map(|s| s.grace_period_millis).unwrap_or_else(default_grace_period_millis); + let grace_period_millis = self.config.reverse_proxy.as_ref().and_then(|r| r.stream.as_ref()).map_or_else(default_grace_period_millis, |s| s.grace_period_millis); self.active_users.connection_permission(username, max_connections, grace_period_millis > 0).await } } diff --git a/src/api/model/hls_cache.rs b/src/api/model/hls_cache.rs index 78ee136ce..4ba60b114 100644 --- a/src/api/model/hls_cache.rs +++ b/src/api/model/hls_cache.rs @@ -70,6 +70,6 @@ impl HlsCache { self.counter.store(1, std::sync::atomic::Ordering::SeqCst); return 1; } - return token; + token } } \ No newline at end of file diff --git a/src/api/model/streams/active_client_stream.rs b/src/api/model/streams/active_client_stream.rs index 904cb51ab..4f0d1b41e 100644 --- a/src/api/model/streams/active_client_stream.rs +++ b/src/api/model/streams/active_client_stream.rs @@ -37,7 +37,7 @@ impl ActiveClientStream { let active_user = app_state.active_users.clone(); let active_provider = app_state.active_provider.clone(); let log_active_clients = app_state.config.log.as_ref().is_some_and(|l| l.active_clients); - let (client_count, connection_count) = active_user.add_connection(&username).await; + let (client_count, connection_count) = active_user.add_connection(username).await; if log_active_clients { info!("Active clients: {client_count}, active connections {connection_count}"); } diff --git a/src/api/model/streams/provider_stream.rs b/src/api/model/streams/provider_stream.rs index a9fae531c..e46b02262 100644 --- a/src/api/model/streams/provider_stream.rs +++ b/src/api/model/streams/provider_stream.rs @@ -26,7 +26,7 @@ pub enum CustomVideoStreamType { fn create_video_stream(video: Option<&Arc>>, headers: &[(String, String)], log_message: &str) -> ProviderStreamResponse { if let Some(video) = video { - trace!("{}", log_message); + trace!("{log_message}"); let mut response_headers: Vec<(String, String)> = headers.iter() .filter(|(key, _)| !(key.eq("content-type") || key.eq("content-length") || key.contains("range"))) .map(|(key, value)| (key.to_string(), value.to_string())).collect(); @@ -76,7 +76,7 @@ pub fn get_header_filter_for_item_type(item_type: PlaylistItemType) -> HeaderFil pub async fn get_provider_pipe_stream(app_state: &AppState, stream_url: &Url, req_headers: &HeaderMap, - input_headers: Option>, + input_headers: Option<&HashMap>, item_type: PlaylistItemType) -> ProviderStreamResponse { let filter_header = get_header_filter_for_item_type(item_type); let req_headers = get_headers_from_request(req_headers, &filter_header); @@ -84,7 +84,7 @@ pub async fn get_provider_pipe_stream(app_state: &AppState, // These are the configured headers for this input. // The stream url, we need to clone it because of move to async block. // We merge configured input headers with the headers from the request. - let headers = get_request_headers(input_headers.as_ref(), Some(&req_headers)); + let headers = get_request_headers(input_headers, Some(&req_headers)); let client = app_state.http_client.get(stream_url.clone()).headers(headers.clone()); match client.send().await { Ok(response) => { @@ -114,7 +114,7 @@ pub async fn get_provider_pipe_stream(app_state: &AppState, pub async fn get_provider_reconnect_buffered_stream(app_state: &AppState, stream_url: &Url, req_headers: &HeaderMap, - input_headers: Option>, + input_headers: Option<&HashMap>, options: BufferStreamOptions) -> ProviderStreamResponse { match create_provider_stream(&app_state.config, Arc::clone(&app_state.http_client), stream_url, req_headers, input_headers, options).await { None => (None, None), diff --git a/src/api/model/streams/provider_stream_factory.rs b/src/api/model/streams/provider_stream_factory.rs index 79adca4e5..518584ee5 100644 --- a/src/api/model/streams/provider_stream_factory.rs +++ b/src/api/model/streams/provider_stream_factory.rs @@ -174,7 +174,7 @@ fn get_request_range_start_bytes(req_headers: &HashMap>) -> Opti fn get_client_stream_request_params( req_headers: &HeaderMap, - input_headers: Option>, + input_headers: Option<&HashMap>, options: &BufferStreamOptions) -> (usize, Option, bool, u32, HeaderMap) { let stream_buffer_size = if options.is_buffer_enabled() { options.get_stream_buffer_size() } else { 0 }; @@ -186,7 +186,7 @@ fn get_client_stream_request_params( req_headers.remove("range"); // We merge configured input headers with the headers from the request. - let headers = get_request_headers(input_headers.as_ref(), Some(&req_headers)); + let headers = get_request_headers(input_headers, Some(&req_headers)); (stream_buffer_size, req_range_start_bytes, options.is_reconnect_enabled(), options.force_reconnect_secs, headers) } @@ -287,7 +287,7 @@ async fn stream_provider(client: Arc, stream_options: ProviderS } } Err(err) => { - debug!("Server connection failed with {err}") + debug!("Server connection failed with {err}"); } } if !stream_options.should_continue() { @@ -317,7 +317,7 @@ async fn get_initial_stream(cfg: &Config, client: Arc, stream_o warn!("The stream could be unavailable. ({status}) {}", sanitize_sensitive_info(stream_options.get_url().as_str())); } } - }; + } if connect_err > ERR_MAX_RETRY_COUNT { break; } @@ -334,7 +334,7 @@ async fn get_initial_stream(cfg: &Config, client: Arc, stream_o fn create_provider_stream_options(stream_url: &Url, req_headers: &HeaderMap, - input_headers: Option>, + input_headers: Option<&HashMap>, options: &BufferStreamOptions) -> ProviderStreamOptions { let (buffer_size, req_range_start_bytes, reconnect, reconnect_force_secs, headers) = get_client_stream_request_params(req_headers, input_headers, options); @@ -356,7 +356,7 @@ pub async fn create_provider_stream(cfg: &Config, client: Arc, stream_url: &Url, req_headers: &HeaderMap, - input_headers: Option>, + input_headers: Option<&HashMap>, options: BufferStreamOptions) -> Option { let stream_options = create_provider_stream_options(stream_url, req_headers, input_headers, &options); diff --git a/src/api/model/streams/throttled_stream.rs b/src/api/model/streams/throttled_stream.rs index f0726134c..dbf0c944f 100644 --- a/src/api/model/streams/throttled_stream.rs +++ b/src/api/model/streams/throttled_stream.rs @@ -16,6 +16,7 @@ pub struct ThrottledStream { } impl ThrottledStream { + #[allow(clippy::cast_precision_loss)] pub fn new(inner: S, throttle_kbps: usize) -> Self { assert!(throttle_kbps > 0, "Rate must be greater than 0"); let rate_bytes_per_sec = (throttle_kbps as f64) * 1000.0 / 8.0; @@ -33,6 +34,7 @@ where { type Item = Result; + #[allow(clippy::cast_precision_loss)] fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { let this = &mut *self; diff --git a/src/auth/access_token.rs b/src/auth/access_token.rs index 5aa70ace8..c2eab8a20 100644 --- a/src/auth/access_token.rs +++ b/src/auth/access_token.rs @@ -1,6 +1,6 @@ -use chrono::{Utc}; -use serde::{Serialize, Deserialize}; -use crate::repository::storage::hex_encode; +use crate::repository::storage::{hex_decode, hex_encode}; +use chrono::Utc; +use serde::{Deserialize, Serialize}; fn constant_time_eq(a: &[u8], b: &[u8]) -> bool { a.len() == b.len() && a.iter().zip(b.iter()).fold(0u8, |acc, (x, y)| acc | (x ^ y)) == 0 @@ -8,50 +8,56 @@ fn constant_time_eq(a: &[u8], b: &[u8]) -> bool { #[derive(Serialize, Deserialize, Debug)] struct AccessToken { - timestamp: i64, - ttl_secs: i64, - signature: String, + ts: i64, + ttl: i64, + sig: String, } -pub fn create_access_token(secret: &[u8; 32], ttl_secs: i64) -> String { +pub fn create_access_token(secret: &[u8; 32], ttl_secs: u16) -> String { let timestamp = Utc::now().timestamp(); - let data = format!("{}", timestamp); - - let hash = blake3::keyed_hash(secret, data.as_bytes()); + let timestamp_bytes = timestamp.to_le_bytes(); + let ttl_secs_bytes = ttl_secs.to_le_bytes(); + let hash = blake3::keyed_hash(secret, ×tamp_bytes); let signature = hex_encode(hash.as_bytes()); - - // token as json - let token = AccessToken { - timestamp, - ttl_secs, - signature, - }; - - // serialize as json - serde_json::to_string(&token).unwrap() + format!("{}{}{signature}", hex_encode(×tamp_bytes), hex_encode(&ttl_secs_bytes)) } pub fn verify_access_token(token_str: &str, secret: &[u8; 32]) -> bool { - // deserialize token - let token: AccessToken = serde_json::from_str(token_str).unwrap(); - - // Validate time - let current_timestamp = Utc::now().timestamp(); - if current_timestamp - token.timestamp > token.ttl_secs { + if token_str.len() < 52 { return false; } - // Create HMAC-Hash for the timestamp with blake3 - let data = token.timestamp.to_string(); - let expected = blake3::keyed_hash(secret, data.as_bytes()); - let expected_hash = hex_encode(expected.as_bytes()); - constant_time_eq(expected_hash.as_bytes(), token.signature.as_bytes()) + let timestamp_bytes = hex_decode(&token_str[0..16]).unwrap_or_default(); + if timestamp_bytes.len() != 8 { + return false; + } + + let timestamp = i64::from_le_bytes(timestamp_bytes.try_into().unwrap_or([0; 8])); + + if timestamp == 0 { + return false; + } + + let ttl_bytes = hex_decode(&token_str[16..20]).unwrap_or_default(); + if ttl_bytes.len() != 2 { + return false; + } + let ttl_secs = u16::from_le_bytes(ttl_bytes.try_into().unwrap_or([0; 2])); + let signature = hex_decode(&token_str[20..]).unwrap_or_default(); + + let current_timestamp = Utc::now().timestamp(); + if current_timestamp - timestamp > i64::from(ttl_secs) { + return false; + } + + let expected = blake3::keyed_hash(secret, ×tamp.to_le_bytes()); + constant_time_eq(expected.as_bytes(), &signature) } #[cfg(test)] mod tests { - use std::thread; use crate::auth::access_token::{create_access_token, verify_access_token}; + use std::thread; #[test] fn test_valid_token() { diff --git a/src/foundation/filter.rs b/src/foundation/filter.rs index 97cede7bc..ef937323c 100644 --- a/src/foundation/filter.rs +++ b/src/foundation/filter.rs @@ -35,7 +35,7 @@ pub fn set_field_value(pli: &mut PlaylistItem, field: &ItemField, value: String) ItemField::Url => header.url = value, ItemField::Input => header.input_name = value, ItemField::Type => {} - }; + } } pub struct ValueProvider<'a> { diff --git a/src/main.rs b/src/main.rs index c5b276235..17aae8e12 100644 --- a/src/main.rs +++ b/src/main.rs @@ -122,11 +122,11 @@ fn main() { let targets = validate_targets(args.target.as_ref(), &cfg.sources).unwrap_or_else(|err| exit!("{}", err)); - info!("Version: {}", VERSION); + info!("Version: {VERSION}"); if let Some(bts) = BUILD_TIMESTAMP.to_string().parse::>().ok().map(|datetime| datetime.format("%Y-%m-%d %H:%M:%S %Z").to_string()) { info!("Build time: {bts}"); } - info!("Current time: {}", chrono::offset::Local::now().format("%Y-%m-%d %H:%M:%S").to_string()); + info!("Current time: {}", chrono::offset::Local::now().format("%Y-%m-%d %H:%M:%S")); info!("Working dir: {:?}", &cfg.working_dir); info!("Config dir: {:?}", &cfg.t_config_path); info!("Config file: {config_file:?}"); @@ -204,7 +204,7 @@ async fn start_in_cli_mode(cfg: Arc, targets: Arc) { async fn start_in_server_mode(cfg: Arc, targets: Arc) { if let Err(err) = api::main_api::start_server(cfg, targets).await { exit!("Can't start server: {err}"); - }; + } } fn get_log_level(log_level: &str) -> LevelFilter { diff --git a/src/messaging.rs b/src/messaging.rs index eb3ae7729..ff76df17f 100644 --- a/src/messaging.rs +++ b/src/messaging.rs @@ -32,7 +32,7 @@ fn send_http_post_request(msg: &str, messaging: &MessagingConfig) { .await { Ok(_) => debug!("Text message sent successfully to rest api"), - Err(e) => error!("Text message wasn't sent to rest api because of: {}", e), + Err(e) => error!("Text message wasn't sent to rest api because of: {e}"), } }); } @@ -43,8 +43,8 @@ fn send_telegram_message(msg: &str, messaging: &MessagingConfig) { for chat_id in &telegram.chat_ids { let bot = rustelebot::create_instance(&telegram.bot_token, chat_id); match rustelebot::send_message(&bot, msg, None) { - Ok(()) => debug!("Text message sent successfully to {}", chat_id), - Err(e) => error!("Text message wasn't sent to {} because of: {}", chat_id, e) + Ok(()) => debug!("Text message sent successfully to {chat_id}"), + Err(e) => error!("Text message wasn't sent to {chat_id} because of: {e}") } } } diff --git a/src/model/api_proxy.rs b/src/model/api_proxy.rs index b2f4a9eba..5a6fc8a80 100644 --- a/src/model/api_proxy.rs +++ b/src/model/api_proxy.rs @@ -340,7 +340,7 @@ impl ApiProxyConfig { println!("{err}"); errors.push(err.to_string()); } - }; + } } else { let user_db_path = get_api_user_db_path(cfg); if user_db_path.exists() { @@ -352,7 +352,7 @@ impl ApiProxyConfig { for stored_credential in &stored_user.credentials { if !target_user.credentials.iter().any(|c| c.username == stored_credential.username) { target_user.credentials.push(stored_credential.clone()); - }; + } } } else { self.user.push(stored_user); @@ -445,7 +445,7 @@ impl ApiProxyConfig { target_user.get_target_name(username, password) { return Some((credentials.clone(), target_name.to_string())); - }; + } } debug!("Could not find any target for user {username}"); None @@ -455,7 +455,7 @@ impl ApiProxyConfig { for target_user in &self.user { if let Some((credentials, target_name)) = target_user.get_target_name_by_token(token) { return Some((credentials.clone(), target_name.to_string())); - }; + } } None } diff --git a/src/model/config.rs b/src/model/config.rs index 2f462fae3..8b2b5b526 100644 --- a/src/model/config.rs +++ b/src/model/config.rs @@ -466,7 +466,7 @@ impl ConfigTarget { TargetType::M3u => { hdhomerun_needs_m3u = true; } TargetType::Xtream => { hdhomerun_needs_xtream = true; } _ => return create_m3u_filter_error_result!(M3uFilterErrorKind::Info, "HdHomeRun output option `use_output` only accepts `m3u` or `xtream` for target: {}", self.name), - }; + } } } } @@ -605,7 +605,7 @@ pub struct InputAffix { pub value: String, } -#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, Sequence, PartialEq, Eq, Default)] +#[derive(Debug, Copy, Clone, serde::Serialize, serde::Deserialize, Sequence, PartialEq, Eq, Default)] pub enum InputType { #[serde(rename = "m3u")] #[default] @@ -832,7 +832,7 @@ impl ConfigInput { InputType::Xtream }; - match csv_read_inputs(&self) { + match csv_read_inputs(self) { Ok(mut batch_aliases) => { if !batch_aliases.is_empty() { batch_aliases.reverse(); @@ -864,7 +864,7 @@ impl ConfigInput { } pub fn get_user_info(&self) -> Option { - InputUserInfo::new(self.input_type.clone(), self.username.as_deref(), self.password.as_deref(), &self.url) + InputUserInfo::new(self.input_type, self.username.as_deref(), self.password.as_deref(), &self.url) } } @@ -1674,7 +1674,7 @@ impl Config { return Err(err); } } - }; + } Ok(()) } @@ -1700,7 +1700,7 @@ impl Config { Err(err) => return Err(err) } } - }; + } Ok(()) } diff --git a/src/model/mapping.rs b/src/model/mapping.rs index 6e1fa5623..14736a864 100644 --- a/src/model/mapping.rs +++ b/src/model/mapping.rs @@ -279,9 +279,9 @@ impl MappingValueProcessor<'_> { fn set_property(&mut self, key: &str, value: &str) { if !self.pli.header.set_field(key, value) { - error!("Cant set unknown field {} to {}", key, value); + error!("Cant set unknown field {key} set to {value}"); } - trace!("Property {} set to {}", key, value); + trace!("Property {key} set to {value}"); } fn apply_attributes(&mut self, captured_names: &HashMap<&str, &str>) { @@ -319,7 +319,7 @@ impl MappingValueProcessor<'_> { if let Some(cap_value) = captures.get(cap.as_str()) { captured_tag_values.push(cap_value); } else { - debug!("Cant find any tag match for {}", tag_capture); + debug!("Cant find any tag match for {tag_capture}"); return None; } } @@ -417,7 +417,7 @@ impl ValueProcessor for MappingValueProcessor<'_> { } else { "" }; - debug!("match {}: {}", capture_name, capture_value); + debug!("match {capture_name}: {capture_value}"); captured_values.insert(capture_name.as_str(), capture_value); } ); @@ -492,7 +492,7 @@ impl MappingDefinition { } Err(err) => return Err(err), } - }; + } for mapping in &mut self.mapping { let template_list = self.templates.as_ref(); let tag_list = self.tags.as_ref(); diff --git a/src/model/playlist.rs b/src/model/playlist.rs index 7994168db..9efe23d88 100644 --- a/src/model/playlist.rs +++ b/src/model/playlist.rs @@ -121,7 +121,7 @@ impl FromStr for PlaylistItemType { "LiveUnknown" => Ok(PlaylistItemType::LiveUnknown), "LiveHls" => Ok(PlaylistItemType::LiveHls), "LiveDash" => Ok(PlaylistItemType::LiveDash), - _ => Err(format!("Invalid PlaylistItemType: {}", s)), + _ => Err(format!("Invalid PlaylistItemType: {s}")), } } } diff --git a/src/model/xtream.rs b/src/model/xtream.rs index 3ec2ad0df..70b4e758f 100644 --- a/src/model/xtream.rs +++ b/src/model/xtream.rs @@ -599,7 +599,7 @@ pub fn xtream_playlistitem_to_document(pli: &XtreamPlaylistItem, url: &str, opti XtreamCluster::Series => { document.insert("series_id".to_string(), stream_id_value); } - }; + } let props = pli.additional_properties.as_ref().and_then(|add_props| serde_json::from_str::>(add_props).ok()); @@ -627,7 +627,7 @@ pub fn xtream_playlistitem_to_document(pli: &XtreamPlaylistItem, url: &str, opti append_mandatory_fields(&mut document, SERIES_STREAM_FIELDS); append_release_date(&mut document); } - }; + } rewrite_doc_urls(resource_url.as_ref(), &mut document, XTREAM_VOD_REWRITE_URL_PROPS, ""); @@ -650,7 +650,7 @@ pub fn rewrite_doc_urls(resource_url: Option<&String>, document: &mut Map {} - }; + } } for &field in fields { if let Some(Value::String(value)) = document.get(field) { diff --git a/src/processing/parser/xmltv.rs b/src/processing/parser/xmltv.rs index 2378f9d2b..729d11211 100644 --- a/src/processing/parser/xmltv.rs +++ b/src/processing/parser/xmltv.rs @@ -45,7 +45,7 @@ impl TVGuide { epg_channel_ids.insert(epg_id.to_string()); } std::collections::hash_map::Entry::Vacant(_entry) => {} - }; + } } if epg_channel_ids.contains(epg_id.as_str()) { children.push(tag); @@ -63,7 +63,7 @@ impl TVGuide { tv_attributes.clone_from(&tag.attributes); } _ => {} - }; + } }; parse_tvguide(&mut reader, &mut filter_tags); diff --git a/src/processing/processor/affix.rs b/src/processing/processor/affix.rs index e68a135a1..867718863 100644 --- a/src/processing/processor/affix.rs +++ b/src/processing/processor/affix.rs @@ -22,7 +22,7 @@ fn validate_and_create_affix_processor(affix: Option<&InputAffix>, is_prefix: bo if (valid_property!(&affix_def.field.as_str(), AFFIX_FIELDS) && !affix_def.value.is_empty()) { return Some(create_affix_processor(affix_def, is_prefix)); } - }; + } None } diff --git a/src/processing/processor/playlist.rs b/src/processing/processor/playlist.rs index 3eb779498..f0ac34ecd 100644 --- a/src/processing/processor/playlist.rs +++ b/src/processing/processor/playlist.rs @@ -196,7 +196,7 @@ fn map_channel(mut channel: PlaylistItem, mapping: &Mapping) -> PlaylistItem { if !mapping.mapper.is_empty() { let header = &channel.header; let channel_name = if mapping.match_as_ascii { unidecode(&header.name) } else { header.name.to_string() }; - if mapping.match_as_ascii && log_enabled!(Level::Trace) { trace!("Decoded {} for matching to {}", &header.name, &channel_name); }; + if mapping.match_as_ascii && log_enabled!(Level::Trace) { trace!("Decoded {} for matching to {}", &header.name, &channel_name); } // let ref_chan = &mut channel; let ref_chan = &mut channel; let mut mock_processor = MockValueProcessor {}; @@ -213,7 +213,7 @@ fn map_channel(mut channel: PlaylistItem, mapping: &Mapping) -> PlaylistItem { _ => { apply_pattern!(&m.t_pattern, &provider, &mut processor); } - }; + } } } channel @@ -348,7 +348,7 @@ async fn process_source(client: Arc, cfg: Arc, source_i } let elapsed = start_time.elapsed().as_secs(); input_stats.insert(input_name.to_string(), create_input_stat(group_count, channel_count, error_list.len(), - input.input_type.clone(), input_name, elapsed)); + input.input_type, input_name, elapsed)); } } if source_playlists.is_empty() { @@ -395,7 +395,7 @@ async fn process_sources(client: Arc, config: Arc, user let thread_num = config.threads; let process_parallel = thread_num > 1 && config.sources.len() > 1; if process_parallel && log_enabled!(Level::Debug) { - debug!("Using {} threads", thread_num); + debug!("Using {thread_num} threads"); } let errors = Arc::new(Mutex::>::new(vec![])); let stats = Arc::new(Mutex::>::new(vec![])); @@ -497,7 +497,7 @@ fn flatten_groups(playlistgroups: Vec) -> Vec { std::collections::hash_map::Entry::Occupied(o) => { sort_order.get_mut(*o.get()).unwrap().channels.extend(group.channels); } - }; + } } sort_order } @@ -543,7 +543,7 @@ async fn process_playlist_for_target(client: Arc, match channel.header.epg_channel_id.as_ref() { None => {normalized_epg_channel_ids.insert(normalize_channel_name(&channel.header.name), None);}, Some(epg_id) => {epg_channel_ids.insert(epg_id.to_string());}, - }; + } } // let epg_channel_ids: HashSet<_> = fp.playlistgroups.iter().flat_map(|g| &g.channels) // .filter_map(|c| c.header.epg_channel_id.as_ref()).map(|a| a.as_str()).collect(); @@ -568,10 +568,10 @@ async fn process_playlist_for_target(client: Arc, if c.header.epg_channel_id.is_some() && (c.header.logo.is_empty() || c.header.logo_small.is_empty()) { if let Some(icon) = epg_icons.get(c.header.epg_channel_id.as_ref().unwrap()) { if c.header.logo.is_empty() { - c.header.logo = icon.to_string(); + c.header.logo = (*icon).to_string(); } if c.header.logo_small.is_empty() { - c.header.logo = icon.to_string(); + c.header.logo = (*icon).to_string(); } } } @@ -619,7 +619,7 @@ pub async fn exec_processing(client: Arc, cfg: Arc, tar } if let Ok(stats_msg) = serde_json::to_string(&serde_json::Value::Object(serde_json::map::Map::from_iter([("stats".to_string(), serde_json::to_value(stats).unwrap())]))) { // print stats - info!("{}", stats_msg); + info!("{stats_msg}"); // send stats send_message(&MsgKind::Stats, cfg.messaging.as_ref(), stats_msg.as_str()); } diff --git a/src/repository/bplustree.rs b/src/repository/bplustree.rs index 85442af6c..8cfaa729e 100644 --- a/src/repository/bplustree.rs +++ b/src/repository/bplustree.rs @@ -534,7 +534,7 @@ where error!("Failed to read id tree from file {err}"); return None; } - }; + } } } // @@ -671,7 +671,7 @@ where error!("Failed to read id tree from file {err}"); return Err(io::Error::new(io::ErrorKind::NotFound, format!("Failed to read id tree from file {err}"))); } - }; + } } } } diff --git a/src/repository/indexed_document.rs b/src/repository/indexed_document.rs index 07da4af7e..c264f4ab1 100644 --- a/src/repository/indexed_document.rs +++ b/src/repository/indexed_document.rs @@ -116,7 +116,7 @@ where // Initialize the index tree (BPlusTree) - either by deserializing an existing one or creating a new one let index_tree = if append_mode && index_path.exists() { IndexedDocumentIndex::::load(&index_path).unwrap_or_else(|err| { - error!("Failed to load index {:?}: {}", index_path, err); + error!("Failed to load index {index_path:?}: {err}"); IndexedDocumentIndex::::new() }) } else { diff --git a/src/repository/kodi_repository.rs b/src/repository/kodi_repository.rs index 1f0ec438b..a405af5b6 100644 --- a/src/repository/kodi_repository.rs +++ b/src/repository/kodi_repository.rs @@ -233,14 +233,14 @@ async fn kodi_style_rename( if let Some(value) = season { filename.push(format!("{separator}S{value:02}")); file_dir.push(format!("Season{separator}{value}")); - }; + } if let Some(value) = episode { if season.is_none() { filename.push(separator.to_string()); } filename.push(format!("E{value:02}")); - }; + } if let Some(value) = title { let sanitized_value = sanitize_for_filename( &trim_whitespace( @@ -342,7 +342,7 @@ async fn get_tmdb_value( _ => {} } }; - }; + } entry.insert(None); None } @@ -563,7 +563,7 @@ async fn prepare_strm_files( let filename = Arc::new(strm_file_name); if all_filenames.contains(&filename) { collisions.insert(Arc::clone(&filename)); - }; + } all_filenames.insert(Arc::clone(&filename)); result.push(StrmFile { file_name: Arc::clone(&filename), @@ -681,7 +681,7 @@ pub async fn kodi_write_strm_playlist( if let Err(err) = write_strm_index_file(cfg, &processed_strm, &strm_index_path).await { failed.push(err); - }; + } if let Err(err) = cleanup_strm_output_directory(target_output.cleanup, &root_path, &existing_strm, &processed_strm).await @@ -730,7 +730,7 @@ async fn ensure_strm_file_directory(failed: &mut Vec, output_path: &Path if let Err(e) = create_dir_all(output_path).await { let err_msg = format!("Failed to create directory for strm playlist: {output_path:?} {e}"); - error!("{}", err_msg); + error!("{err_msg}"); failed.push(err_msg); return false; // skip creation, could not create directory }; @@ -775,7 +775,7 @@ async fn has_strm_file_same_hash(file_path: &PathBuf, content_hash: UUIDType) -> Err(err) => { error!("Could not read existing strm file {file_path:?} {err}"); } - }; + } } false } diff --git a/src/repository/m3u_playlist_iterator.rs b/src/repository/m3u_playlist_iterator.rs index ee204f6b5..3b2d31a50 100644 --- a/src/repository/m3u_playlist_iterator.rs +++ b/src/repository/m3u_playlist_iterator.rs @@ -69,7 +69,7 @@ impl M3uPlaylistIterator { PlaylistItemType::Live | PlaylistItemType::Catchup | PlaylistItemType::LiveUnknown - | PlaylistItemType::LiveHls => "live", + | PlaylistItemType::LiveHls | PlaylistItemType::LiveDash => "live", PlaylistItemType::Video => "movie", PlaylistItemType::Series diff --git a/src/repository/storage.rs b/src/repository/storage.rs index 8ac586f61..0c09c10a0 100644 --- a/src/repository/storage.rs +++ b/src/repository/storage.rs @@ -31,6 +31,20 @@ pub fn hex_encode(bytes: &[u8]) -> String { }) } +pub fn hex_decode(hex: &str) -> Result, String> { + if hex.len() % 2 != 0 { + return Err("hex string must have even length".to_string()); + } + + (0..hex.len()) + .step_by(2) + .map(|i| { + u8::from_str_radix(&hex[i..i+2], 16) + .map_err(|e| format!("invalid hex at position {i}: {e}")) + }) + .collect() +} + pub fn hash_string_as_hex(url: &str) -> String { hex_encode(&hash_string(url)) } diff --git a/src/repository/xtream_repository.rs b/src/repository/xtream_repository.rs index b193c2314..88ca6c1af 100644 --- a/src/repository/xtream_repository.rs +++ b/src/repository/xtream_repository.rs @@ -440,7 +440,7 @@ pub fn xtream_write_series_info( let new_record = record.copy_update_timestamp(); let _ = target_id_mapping.update(&series_info_id, new_record); } - }; + } } Ok(()) @@ -477,7 +477,7 @@ fn xtream_get_info_mapping(config: &Config, target_name: &str, info_id: u32) -> let target_id_mapping_file = get_target_id_mapping_file(&target_path); let _file_lock = config.file_locks.read_lock(&target_id_mapping_file); BPlusTreeQuery::::try_new(&target_id_mapping_file).map_err(|err| { - error!("Could not load id mapping for target {target_name}: {}", err); + error!("Could not load id mapping for target {target_name}: {err}"); str_to_io_error(&format!("ID mapping load error for target {target_name}")) }).ok().map(|mut tree| tree.query(&info_id))? } @@ -500,7 +500,7 @@ pub fn xtream_load_series_info( return match IndexedDocumentDirectAccess::read_indexed_item::(&info_path, &idx_path, &series_id) { Ok(content) => Some(content), Err(err) => { - error!("Failed to read series info for id {series_id} for {target_name}: {}", err); + error!("Failed to read series info for id {series_id} for {target_name}: {err}"); None } }; @@ -534,14 +534,7 @@ pub fn xtream_load_vod_info( if info_path.exists() && idx_path.exists() { { let _file_lock = config.file_locks.read_lock(&info_path); - return match IndexedDocumentDirectAccess::read_indexed_item::(&info_path, &idx_path, &vod_id) { - Ok(content) => Some(content), - Err(_err) => { - // this is not an error, it means the info is not indexed - // error!("Failed to read vod info for id {vod_id} for {target_name}: {}",err); - None - } - }; + return IndexedDocumentDirectAccess::read_indexed_item::(&info_path, &idx_path, &vod_id).ok(); } } None @@ -568,7 +561,7 @@ async fn rewrite_xtream_vod_info

( // doc.insert(TAG_INFO_DATA, Value::Object(info_data)); } ProxyType::Redirect => {} - }; + } } } @@ -698,7 +691,7 @@ async fn rewrite_xtream_series_info

( } if let Err(err) = target_id_mapping.persist() { - error!("{}", err.to_string()); + error!("{err}"); } drop(file_lock); drop(target_id_mapping); diff --git a/src/utils/compression/compression_utils.rs b/src/utils/compression/compression_utils.rs index c5cb7845c..69452192d 100644 --- a/src/utils/compression/compression_utils.rs +++ b/src/utils/compression/compression_utils.rs @@ -1,3 +1,7 @@ +use std::io::{Read, Write}; +use flate2::Compression; +use flate2::read::GzDecoder; +use flate2::write::GzEncoder; pub const ENCODING_GZIP: &str = "gzip"; pub const ENCODING_DEFLATE: &str = "deflate"; @@ -9,4 +13,17 @@ pub const fn is_gzip(bytes: &[u8]) -> bool { pub const fn is_deflate(bytes: &[u8]) -> bool { bytes[0] == 0x78 && (bytes[1] == 0x01 || bytes[1] == 0x9C || bytes[1] == 0xDA) -} \ No newline at end of file +} + +pub fn compress_string(input: &str) -> std::io::Result> { + let mut encoder = GzEncoder::new(Vec::new(), Compression::default()); + encoder.write_all(input.as_bytes())?; + encoder.finish() +} + +pub fn decompress_string(input: &[u8]) -> std::io::Result { + let mut decoder = GzDecoder::new(input); + let mut decompressed = String::new(); + decoder.read_to_string(&mut decompressed)?; + Ok(decompressed) +} diff --git a/src/utils/file/config_reader.rs b/src/utils/file/config_reader.rs index 9f611b798..110e1f3c6 100644 --- a/src/utils/file/config_reader.rs +++ b/src/utils/file/config_reader.rs @@ -126,7 +126,7 @@ pub fn read_api_proxy(config: &Config, api_proxy_file: &str, resolve_env: bool) Ok(mut result) => { match result.prepare(config) { Err(err) => { - exit!("cant read api-proxy-config file: {}", err); + exit!("cant read api-proxy-config file: {err}"); } _ => { Some(result) @@ -134,7 +134,7 @@ pub fn read_api_proxy(config: &Config, api_proxy_file: &str, resolve_env: bool) } } Err(err) => { - error!("cant read api-proxy-config file: {}", err); + error!("cant read api-proxy-config file: {err}"); None } } @@ -192,26 +192,24 @@ const FIELD_PASSWORD: &str = "password"; const FIELD_UNKNOWN: &str = "?"; const DEFAULT_COLUMNS: &[&str] = &[FIELD_URL, FIELD_MAX_CON, FIELD_PRIO, FIELD_NAME, FIELD_USERNAME, FIELD_PASSWORD]; -fn csv_assign_mandatory_fields(alias: &mut ConfigInputAlias, input_type: &InputType) { +fn csv_assign_mandatory_fields(alias: &mut ConfigInputAlias, input_type: InputType) { if !alias.url.is_empty() { match Url::parse(alias.url.as_str()) { Ok(url) => { let (username, password) = get_credentials_from_url(&url); if username.is_none() || password.is_none() { // xtream url - if input_type == &InputType::XtreamBatch { + if input_type == InputType::XtreamBatch { alias.url = url.origin().ascii_serialization().to_string(); - } else if input_type == &InputType::M3uBatch { - if alias.username.is_some() && alias.password.is_some() { - alias.url = format!("{}/get_php?username={}&password={}&type=m3u_plus", - url.origin().ascii_serialization().to_string(), - alias.username.as_deref().unwrap_or("").to_string(), - alias.password.as_deref().unwrap_or("").to_string() - ) - } + } else if input_type == InputType::M3uBatch && alias.username.is_some() && alias.password.is_some() { + alias.url = format!("{}/get_php?username={}&password={}&type=m3u_plus", + url.origin().ascii_serialization(), + alias.username.as_deref().unwrap_or(""), + alias.password.as_deref().unwrap_or("") + ); } } else { - if input_type == &InputType::XtreamBatch { + if input_type == InputType::XtreamBatch { alias.url = url.origin().ascii_serialization().to_string(); } // m3u url @@ -220,12 +218,12 @@ fn csv_assign_mandatory_fields(alias: &mut ConfigInputAlias, input_type: &InputT } if alias.name.is_empty() { - let username = alias.username.as_ref().map(|s| s.as_str()).unwrap_or_default(); + let username = alias.username.as_deref().unwrap_or_default(); let domain: Vec<&str> = url.domain().unwrap_or_default().split('.').collect(); if domain.len() > 1 { alias.name = format!("{}_{username}", domain[domain.len() - 2]); } else { - alias.name = format!("{username}"); + alias.name = username.to_string(); } } } @@ -273,7 +271,7 @@ pub fn csv_read_inputs_from_reader(batch_input_type: InputType, reader: impl Buf let mut result = vec![]; let mut default_columns = vec![]; default_columns.extend_from_slice(DEFAULT_COLUMNS); - for line in reader.lines().into_iter() { + for line in reader.lines() { let line = line?; if line.is_empty() { continue @@ -298,8 +296,8 @@ pub fn csv_read_inputs_from_reader(batch_input_type: InputType, reader: impl Buf let mut config_input = ConfigInputAlias { id: 0, - name: "".to_string(), - url: "".to_string(), + name: String::new(), + url: String::new(), username: None, password: None, priority: 0, @@ -307,13 +305,12 @@ pub fn csv_read_inputs_from_reader(batch_input_type: InputType, reader: impl Buf }; let columns: Vec<&str> = line.split(CSV_SEPARATOR).collect(); - for (&header, &value) in default_columns.iter().zip(columns.iter()).into_iter() { + for (&header, &value) in default_columns.iter().zip(columns.iter()) { if let Err(err) = csv_assign_config_input_column(&mut config_input, header, value) { - error!("Could not parse input line: {} err: {err}", line); - continue; + error!("Could not parse input line: {line} err: {err}"); } } - csv_assign_mandatory_fields(&mut config_input, &input_type); + csv_assign_mandatory_fields(&mut config_input, input_type); result.push(config_input); } Ok(result) @@ -330,7 +327,7 @@ pub fn csv_read_inputs(input: &ConfigInput) -> Result, io: }; match result { Ok(content) => { - return csv_read_inputs_from_reader(input.input_type.clone(), EnvResolvingReader::new(file_reader(Cursor::new(content)))); + return csv_read_inputs_from_reader(input.input_type, EnvResolvingReader::new(file_reader(Cursor::new(content)))); } Err(err) => { return Err(err) diff --git a/src/utils/file/file_utils.rs b/src/utils/file/file_utils.rs index 970526f50..5afd5c278 100644 --- a/src/utils/file/file_utils.rs +++ b/src/utils/file/file_utils.rs @@ -125,10 +125,10 @@ pub fn persist_file(persist_file: Option, text: &str) { let filename = &path_buf.to_str().unwrap_or("?"); match File::create(&path_buf) { Ok(mut file) => match file.write_all(text.as_bytes()) { - Ok(()) => debug!("persisted: {}", filename), - Err(e) => error!("failed to persist file {}, {}", filename, e) + Ok(()) => debug!("persisted: {filename}"), + Err(e) => error!("failed to persist file {filename}, {e}") }, - Err(e) => error!("failed to persist file {}, {}", filename, e) + Err(e) => error!("failed to persist file {filename}, {e}") } } } diff --git a/src/utils/json_utils.rs b/src/utils/json_utils.rs index 6a06b26b6..386582ef2 100644 --- a/src/utils/json_utils.rs +++ b/src/utils/json_utils.rs @@ -169,10 +169,7 @@ where pub fn get_u64_from_serde_value(value: &Value) -> Option { match value { Value::Number(num_val) => num_val.as_u64(), - Value::String(str_val) => match str_val.parse::() { - Ok(val) => Some(val), - Err(_) => None, - }, + Value::String(str_val) => str_val.parse::().ok(), _ => None, } } diff --git a/src/utils/network/epg.rs b/src/utils/network/epg.rs index cbea69a25..65caf2d6a 100644 --- a/src/utils/network/epg.rs +++ b/src/utils/network/epg.rs @@ -11,7 +11,7 @@ pub async fn get_xmltv(client: Arc, _cfg: &Config, input: &Conf match &input.epg_url { None => (None, vec![]), Some(url) => { - debug!("Getting epg file path for url: {}", url); + debug!("Getting epg file path for url: {url}"); let persist_file_path = prepare_file_path(input.persist.as_deref(), working_dir, "") .map(|path| file_utils::add_prefix_to_filename(&path, "epg_", Some("xml"))); diff --git a/src/utils/network/request.rs b/src/utils/network/request.rs index 4f5f3f3e4..a54652225 100644 --- a/src/utils/network/request.rs +++ b/src/utils/network/request.rs @@ -56,7 +56,7 @@ pub async fn get_input_text_content_as_file(client: Arc, input: return create_m3u_filter_error_result!(M3uFilterErrorKind::Notify, "Failed to persist: {} => {}", to_file.to_str().unwrap_or("?"), e); } } - }; + } if filepath.exists() { Some(filepath) @@ -72,8 +72,8 @@ pub async fn get_input_text_content_as_file(client: Arc, input: result.map_or_else(|| { let msg = format!("cant read input url: {}", sanitize_sensitive_info(url_str)); - error!("{}", msg); - create_m3u_filter_error_result!(M3uFilterErrorKind::Notify, "{}", msg) + error!("{msg}"); + create_m3u_filter_error_result!(M3uFilterErrorKind::Notify, "{msg}") }, Ok) } } @@ -103,7 +103,7 @@ pub async fn get_input_text_content(client: Arc, input: &Config return create_m3u_filter_error_result!(M3uFilterErrorKind::Notify, "Failed to persist: {} => {}", to_file.to_str().unwrap_or("?"), e); } } - }; + } match get_local_file_content(&filepath) { Ok(content) => Some(content), @@ -119,8 +119,8 @@ pub async fn get_input_text_content(client: Arc, input: &Config }; result.map_or_else(|| { let msg = format!("cant read input url: {}", sanitize_sensitive_info(url_str)); - error!("{}", msg); - create_m3u_filter_error_result!(M3uFilterErrorKind::Notify, "{}", msg) + error!("{msg}"); + create_m3u_filter_error_result!(M3uFilterErrorKind::Notify, "{msg}") }, Ok) } } @@ -159,7 +159,7 @@ pub fn get_request_headers(defined_headers: Option<&HashMap>, cu if log_enabled!(Level::Trace) { let he: HashMap = headers.iter().map(|(k, v)| (k.to_string(), String::from_utf8_lossy(v.as_bytes()).to_string())).collect(); if !he.is_empty() { - trace!("Request headers {:?}", he); + trace!("Request headers {he:?}"); } } headers @@ -245,7 +245,7 @@ async fn get_remote_content(client: Arc, input: &ConfigInput, u match decoder.read_to_string(&mut decode_buffer) { Ok(_) => {} Err(err) => return Err(str_to_io_error(&format!("failed to decode gzip content {err}"))) - }; + } } ENCODING_DEFLATE => { let mut decoder = ZlibDecoder::new(&bytes[..]); @@ -255,7 +255,7 @@ async fn get_remote_content(client: Arc, input: &ConfigInput, u } } _ => {} - }; + } } if decode_buffer.is_empty() { @@ -490,7 +490,7 @@ pub fn get_credentials_from_url(url: &Url) -> (Option, Option) { } pub fn get_credentials_from_url_str(url_with_credentials: &str) -> (Option, Option) { - if let Ok(url) = Url::parse(&url_with_credentials) { + if let Ok(url) = Url::parse(url_with_credentials) { get_credentials_from_url(&url) } else { (None, None) diff --git a/src/utils/network/xtream.rs b/src/utils/network/xtream.rs index 3f0af6883..3beeaa355 100644 --- a/src/utils/network/xtream.rs +++ b/src/utils/network/xtream.rs @@ -145,7 +145,7 @@ pub async fn get_xtream_playlist(client: Arc, input: &ConfigInp if let Err(err) = request::get_input_json_content(Arc::clone(&client), input, base_url.as_str(), None).await { warn!("Failed to login xtream account {username} {err}"); return (Vec::with_capacity(0), vec![err]); - }; + } let mut playlist_groups: Vec = Vec::with_capacity(128); diff --git a/src/utils/size_utils.rs b/src/utils/size_utils.rs index fd67cfd9b..f81b74ef5 100644 --- a/src/utils/size_utils.rs +++ b/src/utils/size_utils.rs @@ -81,14 +81,14 @@ pub fn parse_to_kbps(input: &str) -> Result { let speed_str = input.trim(); for (unit, multiplier) in units { - if speed_str.ends_with(unit) { - let number_part = speed_str[..speed_str.len() - unit.len()].trim(); + if let Some(speed_unit) = speed_str.strip_suffix(unit) { + let number_part = speed_unit.trim(); let value = u64::from_str(number_part).map_err(|_| format!("Invalid speed: {number_part}"))?; return value.checked_mul(*multiplier).ok_or_else(|| format!("Speed too large: {speed_str}")); } } - u64::from_str(&speed_str).map_err(|_| format!("Invalid speed: {speed_str}, supported units are {}", units.iter().map(|p| p.0).collect::>().join(","))) + u64::from_str(speed_str).map_err(|_| format!("Invalid speed: {speed_str}, supported units are {}", units.iter().map(|p| p.0).collect::>().join(","))) } #[cfg(test)]