From cd76891ce4ac001472264ebeca31abe7b40d17dc Mon Sep 17 00:00:00 2001 From: euzu Date: Tue, 28 Oct 2025 19:52:46 +0100 Subject: [PATCH] Some refactorings - HdHomerun - WebUI GeoIP - Playlist Updates ... --- README.md | 8 +- backend/src/api/api_utils.rs | 2 +- backend/src/api/endpoints/v1_api.rs | 2 +- backend/src/api/endpoints/v1_api_playlist.rs | 8 +- backend/src/api/endpoints/xtream_api.rs | 12 ++- backend/src/api/hdhomerun_proprietary.rs | 76 ++++++++++--------- backend/src/api/hdhomerun_ssdp.rs | 2 +- backend/src/api/main_api.rs | 15 ++-- .../src/api/model/active_provider_manager.rs | 6 +- backend/src/api/model/provider_config.rs | 8 +- .../model/streams/shared_stream_manager.rs | 17 +++-- backend/src/api/serve.rs | 4 +- backend/src/auth/fingerprint.rs | 2 +- backend/src/messaging.rs | 8 +- backend/src/processing/playlist_watch.rs | 17 ++--- backend/src/processing/processor/playlist.rs | 17 +++-- backend/src/repository/xtream_repository.rs | 2 - backend/src/utils/geoip.rs | 6 +- .../src/app/components/config/config_view.rs | 2 +- .../app/components/dashboard/streams_table.rs | 7 +- shared/Cargo.toml | 3 +- shared/src/model/config/geoip.rs | 6 ++ shared/src/model/config/reverse_proxy.rs | 4 + shared/src/model/stream_info.rs | 9 ++- shared/src/utils/string_utils.rs | 10 +-- 25 files changed, 148 insertions(+), 105 deletions(-) diff --git a/README.md b/README.md index 1a8cfc2d3..66fccb337 100644 --- a/README.md +++ b/README.md @@ -218,7 +218,7 @@ Attributes: - `throttle` Allowed units are `KB/s`,`MB/s`,`KiB/s`,`MiB/s`,`kbps`,`mbps`,`Mibps`. Default unit is `kbps` - `grace_period_millis` default set to 300 milliseconds. - `grace_period_timeout_secs` default set to 2 seconds. -- `geopip` is for resolving ip addresses to country names. +- `geopip` is for resolving IP addresses to country names. ##### 1.6.1.1 `retry` If set to `true` on connection loss to provider, the stream will be reconnected. @@ -279,8 +279,8 @@ It has 2 attributes: url: ``` -The `url` is optional and default vaue is: `https://raw.githubusercontent.com/sapics/ip-location-db/refs/heads/main/asn-country/asn-country-ipv4.csv` -The format is csv with 3 columns `range_start,range_end,country_code` +The `url` is optional; default value: `https://raw.githubusercontent.com/sapics/ip-location-db/refs/heads/main/asn-country/asn-country-ipv4.csv` +The format is CSV with 3 columns: `range_start,range_end,country_code`. Example: ```csv @@ -1518,7 +1518,7 @@ user: status: Active ``` -If yu use a reverse proxy in fron of Tuliprox, dont forget to forward +If you use a reverse proxy in front of Tuliprox, don’t forget to forward: - `X-Real-IP` - `X-Forwarded-For` diff --git a/backend/src/api/api_utils.rs b/backend/src/api/api_utils.rs index 1c9cc921b..98437dde0 100644 --- a/backend/src/api/api_utils.rs +++ b/backend/src/api/api_utils.rs @@ -1144,7 +1144,7 @@ async fn fetch_resource_with_retry( build_stream_response(app_state, resource_url, response).await, ); } - // Retry only for 400, 408, 425, 429 and all 5xx statuses + // Retry only for 408, 425, 429 and all 5xx statuses let should_retry = status.is_server_error() || matches!( status, diff --git a/backend/src/api/endpoints/v1_api.rs b/backend/src/api/endpoints/v1_api.rs index 35743c1d9..63e6c28c2 100644 --- a/backend/src/api/endpoints/v1_api.rs +++ b/backend/src/api/endpoints/v1_api.rs @@ -107,7 +107,7 @@ async fn geoip_update(axum::extract::State(app_state): axum::extract::State { error!("Failed to process geoip db: {err}"); - axum::http::StatusCode::NOT_FOUND.into_response() + axum::http::StatusCode::INTERNAL_SERVER_ERROR.into_response() }, _ => { axum::http::StatusCode::INTERNAL_SERVER_ERROR.into_response() diff --git a/backend/src/api/endpoints/v1_api_playlist.rs b/backend/src/api/endpoints/v1_api_playlist.rs index f73642ce0..4a197f1c4 100644 --- a/backend/src/api/endpoints/v1_api_playlist.rs +++ b/backend/src/api/endpoints/v1_api_playlist.rs @@ -60,10 +60,16 @@ async fn playlist_update( let process_targets = app_state.app_config.sources.load().validate_targets(user_targets.as_ref()); match process_targets { Ok(valid_targets) => { + let http_client = Arc::clone(&app_state.http_client.load()); let app_config = Arc::clone(&app_state.app_config); let event_manager = Arc::clone(&app_state.event_manager); let playlist_state = Arc::clone(&app_state.playlists); - playlist::exec_processing(Arc::clone(&app_state.http_client.load()), app_config, Arc::new(valid_targets), Some(event_manager), Some(playlist_state)).await; + let valid_targets = Arc::new(valid_targets); + tokio::spawn({ + async move { + playlist::exec_processing(http_client, app_config, valid_targets, Some(event_manager), Some(playlist_state)).await; + } + }); axum::http::StatusCode::ACCEPTED.into_response() } Err(err) => { diff --git a/backend/src/api/endpoints/xtream_api.rs b/backend/src/api/endpoints/xtream_api.rs index b5cb3ed32..501fc2f18 100644 --- a/backend/src/api/endpoints/xtream_api.rs +++ b/backend/src/api/endpoints/xtream_api.rs @@ -219,7 +219,17 @@ async fn xtream_player_api_stream( stream_req: ApiStreamRequest<'_>, ) -> impl IntoResponse + Send { - debug!("Stream Request {stream_req:?} - {req_headers:?}"); + // if log_enabled!(log::Level::Debug) { + // debug!( + // "Stream request ctx={} user={} stream_id={} action_path={}", + // stream_req.context, + // sanitize_sensitive_info(stream_req.username), + // sanitize_sensitive_info(stream_req.stream_id), + // sanitize_sensitive_info(stream_req.action_path), + // ); + // let message = format!("Client Request headers {req_headers:?}"); + // debug!("{}", sanitize_sensitive_info(&message)); + // } if log_enabled!(log::Level::Debug) { let message = format!("Client Request headers {req_headers:?}"); diff --git a/backend/src/api/hdhomerun_proprietary.rs b/backend/src/api/hdhomerun_proprietary.rs index 0ad2ad00f..1253d8f16 100644 --- a/backend/src/api/hdhomerun_proprietary.rs +++ b/backend/src/api/hdhomerun_proprietary.rs @@ -3,7 +3,7 @@ use crate::model::{AppConfig, HdHomeRunDeviceConfig}; use bytes::{Buf, BufMut, BytesMut}; use log::{error, info, trace}; use std::collections::HashMap; -use std::io::{Cursor, Read}; +use std::io::{Cursor}; use std::net::{Ipv4Addr, SocketAddr}; use std::sync::Arc; use std::time::Duration; @@ -46,18 +46,18 @@ fn write_tlv_u8(buf: &mut BytesMut, tag: u8, value: u8) { buf.put_u8(value); } -fn write_tlv_str(buf: &mut BytesMut, tag: u8, value: &str) { +fn write_tlv_str(buf: &mut bytes::BytesMut, tag: u8, value: &str) { let bytes = value.as_bytes(); + let Some(len) = u16::try_from(bytes.len()).ok() else { + error!("TLV value too long (max 65535 bytes)"); + return; + }; + + // 1 byte type buf.put_u8(tag); - if bytes.len() < 0x80 { - buf.put_u8(u8::try_from(bytes.len()).unwrap_or(0)); - } else { - let len = u16::try_from(bytes.len()).unwrap_or(0); - let byte_first = 0x80 | ((len & 0x7F) as u8); - let byte_second = ((len >> 7) & 0xFF) as u8; - buf.put_u8(byte_first); - buf.put_u8(byte_second); - } + // 2 bytes length (big-endian) + buf.put_u16(len); + // value bytes buf.put_slice(bytes); } @@ -107,39 +107,45 @@ fn build_discover_response(device: &HdHomeRunDeviceConfig, server_host: &str) -> fn parse_tlv(cursor: &mut Cursor<&[u8]>) -> HashMap> { let mut tags = HashMap::new(); - while cursor.position() < cursor.get_ref().len() as u64 { - let mut tag_buf = [0u8; 1]; - if Read::read_exact(cursor, &mut tag_buf).is_err() { - break; - } - let mut len_buf = [0u8; 1]; - if Read::read_exact(cursor, &mut len_buf).is_err() { - break; - } - - let len = if (len_buf[0] & 0x80) == 0 { - len_buf[0] as usize - } else { - let mut second_byte = [0u8; 1]; - if Read::read_exact(cursor, &mut second_byte).is_err() { - break; + loop { + let pos = match usize::try_from(cursor.position()) { + Ok(pos) => pos, + Err(_err) => { + return tags; } - ((second_byte[0] as usize) << 7) + ((len_buf[0] & 0x7F) as usize) }; + if pos >= cursor.get_ref().len() { + break; + } + let mut tag_buf = [0u8; 1]; + if std::io::Read::read_exact(cursor, &mut tag_buf).is_err() { + break; + } + let tag = tag_buf[0]; - if cursor.get_ref().len() < usize::try_from(cursor.position()).unwrap_or(0) + len { + // 2-Byte length (big-endian) + let mut len_buf = [0u8; 2]; + if std::io::Read::read_exact(cursor, &mut len_buf).is_err() { + break; + } + let len = u16::from_be_bytes(len_buf) as usize; + + // Check for incomplete TLV + let remaining = cursor.get_ref().len() as u64 - cursor.position(); + if remaining < len as u64 { + break; // incomplete or invalid TLV record + } + + let mut val_buf = vec![0u8; len]; + if std::io::Read::read_exact(cursor, &mut val_buf).is_err() { break; } - let mut val_buf = vec![0; len]; - if Read::read_exact(cursor, &mut val_buf).is_err() { - break; - } - - tags.insert(tag_buf[0], val_buf); + tags.insert(tag, val_buf); } + tags } diff --git a/backend/src/api/hdhomerun_ssdp.rs b/backend/src/api/hdhomerun_ssdp.rs index 9ac24e85c..c71264f4d 100644 --- a/backend/src/api/hdhomerun_ssdp.rs +++ b/backend/src/api/hdhomerun_ssdp.rs @@ -52,7 +52,7 @@ async fn ssdp_task_loop(socket: UdpSocket, app_config: Arc, server_ho "upnp:rootdevice", "ssdp:all", ]; - if !supported.contains(&st.as_str()) && st != "ssdp:all" { continue; } + if !supported.contains(&st.as_str()) { continue; } // Randomized delay per MX let delay_ms = (fastrand::u64(0..=mx*1000)).min(2000); if delay_ms > 0 { tokio::time::sleep(Duration::from_millis(delay_ms)).await; } diff --git a/backend/src/api/main_api.rs b/backend/src/api/main_api.rs index 90b396cf0..55d187099 100644 --- a/backend/src/api/main_api.rs +++ b/backend/src/api/main_api.rs @@ -70,14 +70,19 @@ fn create_shared_data( let config = app_config.config.load(); let use_geoip = config.is_geoip_enabled(); - let geoip = if use_geoip { + let geoip = if use_geoip { let path = get_geoip_path(&config.working_dir); let _file_lock = app_config.file_locks.read_lock(&path); - let geoip = GeoIp::load(&path).ok(); - if geoip.is_some() { - info!("GeoIp db loaded"); + match GeoIp::load(&path) { + Ok(db) => { + info!("GeoIp db loaded"); + Arc::new(ArcSwapOption::from(Some(Arc::new(db)))) + } + Err(err) => { + error!("Failed to load GeoIp db: {err}"); + Arc::new(ArcSwapOption::from(None)) + } } - Arc::new(ArcSwapOption::from_pointee(geoip)) } else { Arc::new(ArcSwapOption::from(None)) }; diff --git a/backend/src/api/model/active_provider_manager.rs b/backend/src/api/model/active_provider_manager.rs index 82ea42353..674a05a42 100644 --- a/backend/src/api/model/active_provider_manager.rs +++ b/backend/src/api/model/active_provider_manager.rs @@ -40,7 +40,7 @@ impl ProviderConnectionGuard { fn send_release(&self, config: &Arc) { let provider_config = Arc::clone(config); - if let Err(_err) = &self.release_tx.send(Arc::clone(config)) { + if let Err(_err) = &self.release_tx.send(provider_config.clone()) { // Fallback tokio::spawn(async move { provider_config.release().await; @@ -536,8 +536,8 @@ impl ProviderLineupManager { let on_connection_change: ProviderConnectionChangeCallback = Arc::new(move |_name: &str, connections: usize| { let connection_change_sender = connection_change_sender.clone(); let provider_cfg_name = cfg_name.clone(); - if let Err(err) = connection_change_sender.send((provider_cfg_name, connections)) { - error!("Failed to send connection change: {cfg_name}: {connections}, {err}"); + if let Err(err) = connection_change_sender.send((provider_cfg_name.clone(), connections)) { + error!("Failed to send connection change: {provider_cfg_name}: {connections}, {err}"); } }); diff --git a/backend/src/api/model/provider_config.rs b/backend/src/api/model/provider_config.rs index 97a812fac..c87554abe 100644 --- a/backend/src/api/model/provider_config.rs +++ b/backend/src/api/model/provider_config.rs @@ -260,11 +260,11 @@ impl ProviderConfig { pub async fn release(&self) { let mut guard = self.connection.write().await; + if guard.current_connections == 1 || guard.current_connections > self.max_connections { + guard.granted_grace = false; + guard.grace_ts = 0; + } if guard.current_connections > 0 { - if guard.current_connections == 1 && self.max_connections > 1 { - guard.granted_grace = false; - guard.grace_ts = 0; - } modify_connections!(self, guard, -1); } } diff --git a/backend/src/api/model/streams/shared_stream_manager.rs b/backend/src/api/model/streams/shared_stream_manager.rs index 1d02246ac..d1214445d 100644 --- a/backend/src/api/model/streams/shared_stream_manager.rs +++ b/backend/src/api/model/streams/shared_stream_manager.rs @@ -147,7 +147,6 @@ impl SharedStreamState { let mut loop_cnt = 0; loop { - loop_cnt += 1; tokio::select! { biased; @@ -162,16 +161,18 @@ impl SharedStreamState { debug!("Shared stream client send error: {address} {err}"); break; } - if loop_cnt > YIELD_COUNTER { - tokio::task::yield_now().await; - loop_cnt = 0; - } + loop_cnt += 1; + if loop_cnt >= YIELD_COUNTER { + tokio::task::yield_now().await; + loop_cnt = 0; + } } Err(tokio::sync::broadcast::error::RecvError::Lagged(skipped)) => { trace!("Client lagged behind. Skipped {skipped} messages. {address}"); - if loop_cnt > YIELD_COUNTER { - tokio::task::yield_now().await; - loop_cnt = 0; + loop_cnt += 1; + if loop_cnt >= YIELD_COUNTER { + tokio::task::yield_now().await; + loop_cnt = 0; } } Err(_) => break, diff --git a/backend/src/api/serve.rs b/backend/src/api/serve.rs index f7b03398b..7e8625088 100644 --- a/backend/src/api/serve.rs +++ b/backend/src/api/serve.rs @@ -151,8 +151,8 @@ where break; } () = &mut signal_closed => { - if let Err(_err) = connection_release.send(remote_addr.to_string()) { - let addr = remote_addr.to_string(); + let addr = remote_addr.to_string(); + if let Err(_err) = connection_release.send(addr.clone()) { user_manager_clone.remove_connection(&addr).await; } debug!("Connection gracefully closed: {remote_addr}"); diff --git a/backend/src/auth/fingerprint.rs b/backend/src/auth/fingerprint.rs index 75343d718..de0b6a6af 100644 --- a/backend/src/auth/fingerprint.rs +++ b/backend/src/auth/fingerprint.rs @@ -67,7 +67,7 @@ impl Fingerprint { let client_ip_port =format!("{client_ip}:{}", addr.port()); let ua = user_agent.unwrap_or_else(String::new); - let key = format!("{client_ip }{ua}"); + let key = format!("{client_ip }|{ua}"); Ok(Fingerprint(key, client_ip_port)) } diff --git a/backend/src/messaging.rs b/backend/src/messaging.rs index 29eed38d2..8f7d6d610 100644 --- a/backend/src/messaging.rs +++ b/backend/src/messaging.rs @@ -78,9 +78,11 @@ async fn send_pushover_message(client: &Arc, msg: &str, messagi async fn dispatch_send_message(client: &Arc, kind: MsgKind, cfg: Option<&MessagingConfig>, msg: &str, json: bool) { if let Some(messaging) = cfg { if is_enabled(kind, messaging) { - send_telegram_message(client, msg, messaging, json).await; - send_http_post_request(client, msg, messaging).await; - send_pushover_message(client, msg, messaging).await; + tokio::join!( + send_telegram_message(client, msg, messaging, json), + send_http_post_request(client, msg, messaging), + send_pushover_message(client, msg, messaging) + ); } } } diff --git a/backend/src/processing/playlist_watch.rs b/backend/src/processing/playlist_watch.rs index 3c77bbe22..d619f7f5c 100644 --- a/backend/src/processing/playlist_watch.rs +++ b/backend/src/processing/playlist_watch.rs @@ -22,7 +22,7 @@ pub async fn process_group_watch(client: &Arc, cfg: &Config, ta let save_path = path.as_path(); let mut changed = false; if path.exists() { - if let Some(loaded_tree) = load_watch_tree(&path) { + if let Some(loaded_tree) = load_watch_tree(&path).await { // Find elements in set2 but not in set1 let added_difference: BTreeSet = new_tree.difference(&loaded_tree).cloned().collect(); let removed_difference: BTreeSet = loaded_tree.difference(&new_tree).cloned().collect(); @@ -38,7 +38,7 @@ pub async fn process_group_watch(client: &Arc, cfg: &Config, ta changed = true; } if changed { - match save_watch_tree(save_path, &new_tree) { + match save_watch_tree(save_path, &new_tree).await { Ok(()) => {} Err(err) => { error!("failed to write watch_file {}: {}", save_path.to_str().unwrap_or_default(), err); @@ -64,6 +64,7 @@ async fn handle_watch_notification(client: &Arc, cfg: &Config, let added = added.iter().map(std::string::ToString::to_string).collect::>(); let removed = removed.iter().map(std::string::ToString::to_string).collect::>(); if !added.is_empty() || !removed.is_empty() { + let changes = WatchChanges { target: target_name.to_string(), group: group_name.to_string(), @@ -77,15 +78,13 @@ async fn handle_watch_notification(client: &Arc, cfg: &Config, } } -fn load_watch_tree(path: &Path) -> Option> { - std::fs::read(path).map_or(None, |encoded| { - let decoded = bincode_deserialize(&encoded[..]).ok()?; - Some(decoded) - }) +async fn load_watch_tree(path: &Path) -> Option> { + let encoded = tokio::fs::read(path).await.ok()?; + bincode_deserialize(&encoded[..]).ok() } -fn save_watch_tree(path: &Path, tree: &BTreeSet) -> std::io::Result<()> { +async fn save_watch_tree(path: &Path, tree: &BTreeSet) -> std::io::Result<()> { let encoded: Vec = bincode_serialize(&tree)?; - std::fs::write(path, encoded) + tokio::fs::write(path, encoded).await } diff --git a/backend/src/processing/processor/playlist.rs b/backend/src/processing/processor/playlist.rs index d2a573de7..0f8ae8be4 100644 --- a/backend/src/processing/processor/playlist.rs +++ b/backend/src/processing/processor/playlist.rs @@ -620,17 +620,22 @@ async fn process_epg(processed_fetched_playlists: &mut Vec>) (new_epg, new_playlist) } -async fn process_watch(cfg: &Config, client: &Arc, target: &ConfigTarget, new_playlist: &Vec) -> bool { +async fn process_watch(cfg: &Config, client: &Arc, target: &ConfigTarget, new_playlist: &[PlaylistGroup]) -> bool { if let Some(watches) = &target.watch { if default_as_default().eq_ignore_ascii_case(&target.name) { error!("cant watch a target with no unique name"); - } else { - for pl in new_playlist { - if watches.iter().any(|r| r.is_match(&pl.title)) { - process_group_watch(client, cfg, &target.name, pl).await; - } + return false; + } + + let mut futs = Vec::new(); + for pl in new_playlist { + if watches.iter().any(|r| r.is_match(&pl.title)) { + futs.push(process_group_watch(client, cfg, &target.name, pl)); } } + if !futs.is_empty() { + futures::future::join_all(futs).await; + } true } else { false diff --git a/backend/src/repository/xtream_repository.rs b/backend/src/repository/xtream_repository.rs index 878ae1e7f..c451445e2 100644 --- a/backend/src/repository/xtream_repository.rs +++ b/backend/src/repository/xtream_repository.rs @@ -440,8 +440,6 @@ pub async fn xtream_get_item_for_stream_id( PlaylistItemType::Series => { if let Ok(mut item) = xtream_read_series_item_for_stream_id(app_config, mapping.parent_virtual_id, &storage_path) { item.provider_id = mapping.provider_id; - - Ok(item) } else { xtream_read_item_for_stream_id(app_config, virtual_id, &storage_path, XtreamCluster::Series) diff --git a/backend/src/utils/geoip.rs b/backend/src/utils/geoip.rs index 1eb7bab6c..f5c6fa762 100644 --- a/backend/src/utils/geoip.rs +++ b/backend/src/utils/geoip.rs @@ -29,7 +29,11 @@ impl GeoIp { pub fn import_ipv4_from_csv(&mut self, mut reader: impl BufRead, db_path: &Path) -> std::io::Result { let mut buf = String::new(); - while reader.read_line(&mut buf)? > 0 { + loop { + buf.clear(); + if reader.read_line(&mut buf)? == 0 { break; } + let line = buf.trim(); + if line.is_empty() || line.starts_with('#') { continue; } let line = buf.trim(); if line.is_empty() || line.starts_with('#') { continue; } diff --git a/frontend/src/app/components/config/config_view.rs b/frontend/src/app/components/config/config_view.rs index cf82fdabd..ea5801eb4 100644 --- a/frontend/src/app/components/config/config_view.rs +++ b/frontend/src/app/components/config/config_view.rs @@ -243,7 +243,7 @@ pub fn ConfigView() -> Html {

{ translate.t(LABEL_CONFIG) }

{html_if!(config_ctx.config.is_some_and(|c| c.config.is_geoip_enabled()), { - diff --git a/frontend/src/app/components/dashboard/streams_table.rs b/frontend/src/app/components/dashboard/streams_table.rs index 1159f7fc5..1eb037ea4 100644 --- a/frontend/src/app/components/dashboard/streams_table.rs +++ b/frontend/src/app/components/dashboard/streams_table.rs @@ -30,7 +30,7 @@ const HEADERS: [&str; 12] = [ "LABEL.DURATION" ]; -pub fn format_duration(seconds: u64) -> String { +fn format_duration(seconds: u64) -> String { let hours = seconds / 3600; let minutes = (seconds % 3600) / 60; let seconds = seconds % 60; @@ -68,9 +68,8 @@ pub fn StreamsTable(props: &StreamsTableProps) -> Html { use_effect_with((), move |_| { - Interval::new(1000, || { - update_timestamps(); - }).forget(); + let interval = Interval::new(1000, update_timestamps); + move || drop(interval) }); diff --git a/shared/Cargo.toml b/shared/Cargo.toml index c1abca4ce..9b4a4cee0 100644 --- a/shared/Cargo.toml +++ b/shared/Cargo.toml @@ -22,5 +22,6 @@ zeroize = "1" chrono = "0.4.42" bytes = "1" ciborium = "0.2.2" -#[cfg(target_arch = "wasm32")] + +[target.'cfg(target_arch = "wasm32")'.dependencies] js-sys = "0.3.81" diff --git a/shared/src/model/config/geoip.rs b/shared/src/model/config/geoip.rs index 1acf37b55..136466cbd 100644 --- a/shared/src/model/config/geoip.rs +++ b/shared/src/model/config/geoip.rs @@ -9,3 +9,9 @@ pub struct GeoIpConfigDto { #[serde(default = "default_geoip_url")] pub url: String, } + +impl GeoIpConfigDto { + pub fn is_empty(&self) -> bool { + !self.enabled && self.url.trim().is_empty() + } +} \ No newline at end of file diff --git a/shared/src/model/config/reverse_proxy.rs b/shared/src/model/config/reverse_proxy.rs index f188762d2..531b85e8b 100644 --- a/shared/src/model/config/reverse_proxy.rs +++ b/shared/src/model/config/reverse_proxy.rs @@ -27,6 +27,7 @@ impl ReverseProxyConfigDto { && (self.stream.is_none() || self.stream.as_ref().is_some_and(|s| s.is_empty())) && (self.cache.is_none() || self.cache.as_ref().is_some_and(|c| c.is_empty())) && (self.rate_limit.is_none() || self.rate_limit.as_ref().is_some_and(|r| r.is_empty())) + && (self.geoip.is_none() || self.geoip.as_ref().is_some_and(|g| g.is_empty())) } pub fn clean(&mut self) { @@ -39,6 +40,9 @@ impl ReverseProxyConfigDto { if self.rate_limit.as_ref().is_some_and(|s| s.is_empty()) { self.rate_limit = None; } + if self.geoip.as_ref().is_some_and(|g| g.is_empty()) { + self.geoip = None; + } } pub(crate) fn prepare(&mut self, working_dir: &str) -> Result<(), TuliproxError> { diff --git a/shared/src/model/stream_info.rs b/shared/src/model/stream_info.rs index 80ad077a1..e5e1f38f6 100644 --- a/shared/src/model/stream_info.rs +++ b/shared/src/model/stream_info.rs @@ -1,6 +1,6 @@ use serde::{Deserialize, Serialize}; use crate::model::{M3uPlaylistItem, PlaylistEntry, PlaylistItemType, XtreamCluster, XtreamPlaylistItem}; -use crate::utils::{current_time_secs, StringExt}; +use crate::utils::{current_time_secs, longest}; #[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] pub struct StreamChannel { @@ -21,7 +21,7 @@ impl XtreamPlaylistItem { item_type: self.item_type, cluster: self.xtream_cluster, group: self.group.clone(), - title: String::longest(self.title.as_str(), self.name.as_str()).to_string(), + title: longest(self.title.as_str(), self.name.as_str()).to_string(), url: self.url.clone(), shared: false, } @@ -36,7 +36,7 @@ impl M3uPlaylistItem { item_type: self.item_type, cluster: XtreamCluster::try_from(self.item_type).unwrap_or(XtreamCluster::Live), group: self.group.clone(), - title: String::longest(self.title.as_str(), self.name.as_str()).to_string(), + title: longest(self.title.as_str(), self.name.as_str()).to_string(), url: self.url.clone(), shared: false, } @@ -49,8 +49,11 @@ pub struct StreamInfo { pub channel: StreamChannel, pub provider: String, pub addr: String, + #[serde(default)] pub user_agent: String, + #[serde(default)] pub ts: u64, + #[serde(default, skip_serializing_if = "Option::is_none")] pub country: Option, } diff --git a/shared/src/utils/string_utils.rs b/shared/src/utils/string_utils.rs index 67a957a80..08d18ffa4 100644 --- a/shared/src/utils/string_utils.rs +++ b/shared/src/utils/string_utils.rs @@ -127,14 +127,8 @@ pub fn humanize_snake_case(s: &str) -> String { result } -pub trait StringExt { - fn longest<'a>(a: &'a str, b: &'a str) -> &'a str; -} - -impl StringExt for String { - fn longest<'a>(a: &'a str, b: &'a str) -> &'a str { - if a.len() >= b.len() { a } else { b } - } +pub fn longest<'a>(a: &'a str, b: &'a str) -> &'a str { + if a.len() >= b.len() { a } else { b } } #[cfg(test)]