From 21006dbf0f2a8a437d51512dbf2f58063fc2b3f9 Mon Sep 17 00:00:00 2001 From: euzu Date: Wed, 26 Nov 2025 14:13:31 +0100 Subject: [PATCH] Refactored connection handling to avoid race conditions --- README.md | 2 +- .../src/api/model/active_provider_manager.rs | 38 ++++++++++++------- .../src/api/model/provider_lineup_manager.rs | 1 + backend/src/model/xtream.rs | 1 + backend/src/utils/file/csv_input_reader.rs | 5 ++- frontend/src/services/user_api_service.rs | 2 +- shared/src/model/config/api_user.rs | 2 +- shared/src/utils/serde_utils.rs | 2 + 8 files changed, 36 insertions(+), 17 deletions(-) diff --git a/README.md b/README.md index 9a96ba764..4c324d741 100644 --- a/README.md +++ b/README.md @@ -684,7 +684,7 @@ Each input has the following attributes: - `headers` is optional - `method` can be `GET` or `POST` - `username` only mandatory for type `xtream` -- `pasword` only mandatory for type `xtream` +- `password` only mandatory for type `xtream` - `exp_date` optional, i a date as "YYYY-MM-DD HH:MM:SS" format like `2028-11-30 12:34:12` or Unix timestamp (seconds since epoch) - `options` is optional, + `xtream_skip_live` true or false, live section can be skipped. diff --git a/backend/src/api/model/active_provider_manager.rs b/backend/src/api/model/active_provider_manager.rs index 2ecc48626..12c7bfc7b 100644 --- a/backend/src/api/model/active_provider_manager.rs +++ b/backend/src/api/model/active_provider_manager.rs @@ -134,7 +134,10 @@ impl ActiveProviderManager { }; if let Some(allocation) = single_allocation { - debug!("Released provider connection {:?} for {addr}", allocation.get_provider_name().unwrap_or_default()); + debug!( + "Released provider connection {:?} for {addr}", + allocation.get_provider_name().unwrap_or_default() + ); allocation.release().await; return; } @@ -148,29 +151,38 @@ impl ActiveProviderManager { None => return, // no shared connection }; - if let Some(shared) = connections.shared.by_key.get_mut(&key) { - shared.connections.remove(addr); - if shared.connections.is_empty() { - let allocation_clone = shared.allocation.clone(); - connections.shared.key_by_addr.retain(|_, v| v != &key); - connections.shared.by_key.remove(&key); - Some(allocation_clone) - } else { - None - } + // Clone the SharedAllocation to avoid double mutable borrow + let mut shared = match connections.shared.by_key.get(&key) { + Some(s) => s.clone(), + None => return, + }; + + // Remove this address from the shared connection set + shared.connections.remove(addr); + // Always remove stale key-by-addr entry + connections.shared.key_by_addr.remove(addr); + + if shared.connections.is_empty() { + // If this was the last user of the shared allocation: + connections.shared.by_key.remove(&key); + Some(shared.allocation) } else { + // Update the entry back with the remaining connections + connections.shared.by_key.insert(key, shared); None } }; // release allocation if let Some(allocation) = shared_allocation { - debug!("Released last shared connection for provider {}, releasing allocation {addr}", allocation.get_provider_name().unwrap_or_default()); allocation.release().await; + debug!( + "Released last shared connection for provider {}, releasing allocation {addr}", + allocation.get_provider_name().unwrap_or_default() + ); } } - pub async fn release_handle(&self, handle: &ProviderHandle) { self.release_connection(&handle.client_id).await; } diff --git a/backend/src/api/model/provider_lineup_manager.rs b/backend/src/api/model/provider_lineup_manager.rs index 63ce51f08..9ccd4ae08 100644 --- a/backend/src/api/model/provider_lineup_manager.rs +++ b/backend/src/api/model/provider_lineup_manager.rs @@ -578,6 +578,7 @@ impl ProviderLineupManager { || a_alias.username != b_alias.username || a_alias.password != b_alias.password || a_alias.url != b_alias.url + || a_alias.exp_date != b_alias.exp_date { return true; } diff --git a/backend/src/model/xtream.rs b/backend/src/model/xtream.rs index 3d68ef13b..f740b2eb9 100644 --- a/backend/src/model/xtream.rs +++ b/backend/src/model/xtream.rs @@ -8,6 +8,7 @@ use shared::utils::{deserialize_as_option_string, deserialize_as_string, deseria get_non_empty_str, opt_string_or_number_u32, string_default_on_null, string_or_number_f64, string_or_number_u32}; use std::iter::FromIterator; +#[derive(Debug, Default)] pub struct XtreamLoginInfo { pub status: Option, pub exp_date: Option, diff --git a/backend/src/utils/file/csv_input_reader.rs b/backend/src/utils/file/csv_input_reader.rs index d614470e0..b64e52f12 100644 --- a/backend/src/utils/file/csv_input_reader.rs +++ b/backend/src/utils/file/csv_input_reader.rs @@ -88,7 +88,10 @@ fn csv_assign_config_input_column(config_input: &mut ConfigInputAliasDto, header config_input.password = Some(value.to_string()); } FIELD_EXP_DATE => { - config_input.exp_date = parse_timestamp(value).ok().flatten(); + config_input.exp_date = parse_timestamp(value).unwrap_or_else(|e| { + error!("Failed to parse exp_date '{value}': {e}"); + None + }); } _ => {} } diff --git a/frontend/src/services/user_api_service.rs b/frontend/src/services/user_api_service.rs index 61d93c12f..602420c13 100644 --- a/frontend/src/services/user_api_service.rs +++ b/frontend/src/services/user_api_service.rs @@ -33,7 +33,7 @@ impl UserApiService { } pub async fn save_playlist_bouquet(&self, bouquet: &PlaylistBouquetDto) -> Result<(), Error> { - match request_post::<&PlaylistBouquetDto, ()>(&self.user_playlist_bouquet_path, bouquet, None, None).await{ + match request_post::<&PlaylistBouquetDto, ()>(&self.user_playlist_bouquet_path, bouquet, None, None).await { Ok(_) => { Ok(()) }, Err(err) => { error!("{err}"); diff --git a/shared/src/model/config/api_user.rs b/shared/src/model/config/api_user.rs index 9beb44512..fa79dc998 100644 --- a/shared/src/model/config/api_user.rs +++ b/shared/src/model/config/api_user.rs @@ -72,7 +72,7 @@ impl ProxyUserCredentialsDto { } } if let Some(exp_date) = self.exp_date { - let now = chrono::Local::now(); + let now = chrono::Utc::now(); if (exp_date - now.timestamp()) < 0 { return false; } diff --git a/shared/src/utils/serde_utils.rs b/shared/src/utils/serde_utils.rs index e70eadea9..3606cc989 100644 --- a/shared/src/utils/serde_utils.rs +++ b/shared/src/utils/serde_utils.rs @@ -148,6 +148,8 @@ where serializer.serialize_str(&u8_16_to_hex(bytes)) } +/// Deserializes a timestamp from either a Unix timestamp (seconds) or a UTC datetime string +/// in the format "YYYY-MM-DD HH:MM:SS". Note: Datetime strings are interpreted as UTC. pub fn deserialize_timestamp<'de, D>(deserializer: D) -> Result, D::Error> where D: serde::Deserializer<'de>,