mirror of
https://github.com/euzu/tuliprox.git
synced 2026-10-07 08:22:05 +02:00
Refactored connection handling to avoid race conditions
This commit is contained in:
@@ -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.
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
@@ -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<ProxyUserStatus>,
|
||||
pub exp_date: Option<i64>,
|
||||
|
||||
@@ -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
|
||||
});
|
||||
}
|
||||
_ => {}
|
||||
}
|
||||
|
||||
@@ -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}");
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
@@ -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<Option<i64>, D::Error>
|
||||
where
|
||||
D: serde::Deserializer<'de>,
|
||||
|
||||
Reference in New Issue
Block a user