diff --git a/CHANGELOG.md b/CHANGELOG.md index a0132375f..024de4242 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -23,7 +23,11 @@ - Added `order: none` support for group/channel sorting so mappings can opt out of any reordering and keep the source order. - Session tracking now matches repeated HLS segment connections by session token so a single user keeps one active connection count even when new TCP sockets are opened. - EPG icon urls are now rewritten on reverse proxy mode. -- Xtream Codes Batch provider accounts are now checked for expiration. +- Short EPG is now served from local disk, if available +- WebUI Api-User Category selection implemented +- Stream Table Copy-To-Clipboard functions added +- Refactored provider connection handling to avoid possible race conditions +- Added `exp_date` field to inputs, aliases, and CSV batch files; accepts date in "YYYY-MM-DD HH:MM:SS" format or Unix timestamp (seconds since epoch). # 3.2.0 (2025-11-14) - Added `name` attribute to Staged Input. diff --git a/Cargo.lock b/Cargo.lock index 3ad5c8441..3ce3606a5 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1096,7 +1096,7 @@ dependencies = [ [[package]] name = "frontend" -version = "3.2.10" +version = "3.2.11" dependencies = [ "anyhow", "base64", @@ -3765,7 +3765,7 @@ dependencies = [ [[package]] name = "shared" -version = "3.2.10" +version = "3.2.11" dependencies = [ "base64", "bitflags 2.10.0", @@ -4314,7 +4314,7 @@ checksum = "e421abadd41a4225275504ea4d6566923418b7f05506fbc9c0fe86ba7396114b" [[package]] name = "tuliprox" -version = "3.2.10" +version = "3.2.11" dependencies = [ "arc-swap", "async-compression", diff --git a/README.md b/README.md index f40345fae..4c324d741 100644 --- a/README.md +++ b/README.md @@ -665,9 +665,6 @@ The template can now be used for sequence - '(?i)\bSD\b' ``` - - - ### 2.2. `sources` `sources` is a sequence of source definitions, which have two top level entries: -`inputs` @@ -687,7 +684,8 @@ 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. + `xtream_skip_vod` true or false, vod section can be skipped. @@ -834,9 +832,9 @@ There are 2 batch input types `xtream_batch` and `m3u_batch`. ``` ```csv -#name;username;password;url;max_connections;priority -my_provider_1;user1;password1;http://my_provider_1.com:80;1;0 -my_provider_2;user2;password2;http://my_provider_2.com:8080;1;0 +#name;username;password;url;max_connections;priority;exp_date +my_provider_1;user1;password1;http://my_provider_1.com:80;1;0;2028-11-23 12:34:23 +my_provider_2;user2;password2;http://my_provider_2.com:8080;1;0;2028-11-23 12:34:23 ``` ##### `M3uBatch` @@ -864,6 +862,11 @@ A `priority` of `0` is higher than `1` Higher numbers mean **lower priority** This means tasks or items with smaller (even negative) values will be handled before those with larger values. +The `exp_date` field is a date as: +- "YYYY-MM-DD HH:MM:SS" format like `2028-11-30 12:34:12` +- or Unix timestamp (seconds since epoch) + + ### 2.2.2 `targets` Has the following top level entries: - `enabled` _optional_ default is `true`, if you disable the processing is skipped diff --git a/backend/Cargo.toml b/backend/Cargo.toml index 39a1f4f12..2b672bde0 100644 --- a/backend/Cargo.toml +++ b/backend/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "tuliprox" -version = "3.2.10" +version = "3.2.11" edition = "2021" # See more keys and their definitions at https://doc.rust-lang.org/cargo/reference/manifest.html diff --git a/backend/src/api/api_utils.rs b/backend/src/api/api_utils.rs index 107ca29de..c683811eb 100644 --- a/backend/src/api/api_utils.rs +++ b/backend/src/api/api_utils.rs @@ -360,7 +360,7 @@ async fn resolve_streaming_strategy( if let Some(allocation) = provider_connection_handle.as_ref().map(|ph| &ph.allocation) { match allocation { ProviderAllocation::Exhausted => { - debug!("Input {} is exhausted. No connections allowed.", input.name); + debug!("Provider {} is exhausted. No connections allowed.", input.name); let stream = create_provider_connections_exhausted_stream(&app_state.app_config, &[]); ProviderStreamState::Custom(stream) } @@ -391,7 +391,7 @@ async fn resolve_streaming_strategy( } } } else { - debug!("Input {} is exhausted. No connections allowed.", input.name); + debug!("Provider {} is exhausted. No connections allowed.", input.name); let stream = create_provider_connections_exhausted_stream(&app_state.app_config, &[]); ProviderStreamState::Custom(stream) }; diff --git a/backend/src/api/endpoints/hls_api.rs b/backend/src/api/endpoints/hls_api.rs index 6aac833fe..4cb9c9f77 100644 --- a/backend/src/api/endpoints/hls_api.rs +++ b/backend/src/api/endpoints/hls_api.rs @@ -166,10 +166,11 @@ pub(in crate::api) async fn handle_hls_stream_request( async fn get_stream_channel(app_state: &Arc, target: &Arc, virtual_id: u32) -> Option { if target.has_output(TargetType::Xtream) { if let Ok((pli, _)) = xtream_repository::xtream_get_item_for_stream_id(virtual_id, app_state, target, None).await { - return Some(pli.to_stream_channel()); + return Some(pli.to_stream_channel(target.id)); } } - m3u_get_item_for_stream_id(virtual_id, app_state, target).await.ok().map(|pli| pli.to_stream_channel()) + let target_id = target.id; + m3u_get_item_for_stream_id(virtual_id, app_state, target).await.ok().map(|pli| pli.to_stream_channel(target_id)) } async fn resolve_stream_channel( @@ -181,6 +182,7 @@ async fn resolve_stream_channel( let mut channel = match get_stream_channel(app_state, target, virtual_id).await { Some(channel) => channel, None => StreamChannel { + target_id: target.id, virtual_id, provider_id: 0, item_type: PlaylistItemType::LiveHls, @@ -225,7 +227,7 @@ async fn hls_api_stream( app_state.app_config.get_input_by_id(params.input_id), true, format!( - "Cant find input for target {target_name}, stream_id {virtual_id}, hls" + "Cant find input {} for target {target_name}, stream_id {virtual_id}, hls", params.input_id ) ); diff --git a/backend/src/api/endpoints/m3u_api.rs b/backend/src/api/endpoints/m3u_api.rs index 100c764ea..26e203103 100644 --- a/backend/src/api/endpoints/m3u_api.rs +++ b/backend/src/api/endpoints/m3u_api.rs @@ -119,7 +119,7 @@ async fn m3u_api_stream( .app_config .get_input_by_name(pli.input_name.as_str()), true, - format!("Cant find input for target {target_name}, stream_id {virtual_id}") + format!("Cant find input {} for target {target_name}, stream_id {virtual_id}", pli.input_name) ); let cluster = XtreamCluster::try_from(pli.item_type).unwrap_or(XtreamCluster::Live); @@ -154,7 +154,7 @@ async fn m3u_api_stream( fingerprint, app_state, session, - pli.to_stream_channel(), + pli.to_stream_channel(target.id), req_headers, &input, &user, @@ -224,7 +224,7 @@ async fn m3u_api_stream( fingerprint, app_state, &session_key, - pli.to_stream_channel(), + pli.to_stream_channel(target.id), session_url, req_headers, &input, diff --git a/backend/src/api/endpoints/xmltv_api.rs b/backend/src/api/endpoints/xmltv_api.rs index 90c6c6726..2599281a2 100644 --- a/backend/src/api/endpoints/xmltv_api.rs +++ b/backend/src/api/endpoints/xmltv_api.rs @@ -145,11 +145,12 @@ fn parse_timeshift(time_shift: Option<&String>) -> Option { }) } -async fn serve_epg( +pub async fn serve_epg( app_state: &Arc, epg_path: &Path, user: &ProxyUserCredentials, target: &Arc, + filter: Option, ) -> axum::response::Response { if let Ok(exists) = tokio::fs::try_exists(epg_path).await { if exists { @@ -165,12 +166,12 @@ async fn serve_epg( // Use 0 for timeshift if None let timeshift = parse_timeshift(user.epg_timeshift.as_ref()).unwrap_or(0); - return if timeshift != 0 || rewrite_urls { + return if timeshift != 0 || rewrite_urls || filter.is_some() { let server_info = app_state.app_config.get_user_server_info(user); let base_url = format!("{}/{}/{}/{}/", server_info.get_base_url(), storage_const::EPG_RESOURCE_PATH, &user.username, &user.password); - // Apply timeshift and/or rewrite URLs - serve_epg_with_rewrites(epg_path, timeshift, rewrite_urls, &encrypt_secret, &base_url).await + // Apply timeshift and/or rewrite URLs and/or filter + serve_epg_with_rewrites(epg_path, timeshift, rewrite_urls, &encrypt_secret, &base_url, filter).await } else { // Neither timeshift nor rewrite needed, serve original file serve_file(epg_path, mime::TEXT_XML).await.into_response() @@ -187,6 +188,7 @@ async fn serve_epg_with_rewrites( rewrite_urls: bool, secret: &[u8; 16], base_url: &str, + filter: Option, ) -> axum::response::Response { match tokio::fs::try_exists(epg_path).await { Ok(exists) => { @@ -231,10 +233,75 @@ async fn serve_epg_with_rewrites( let mut buf = Vec::with_capacity(4096); let duration = Duration::minutes(i64::from(offset_minutes)); + let mut skip_depth = None; loop { - match xml_reader.read_event_into_async(&mut buf).await { - Ok(Event::Start(ref e)) if offset_minutes != 0 && e.name().as_ref() == b"programme" => { + buf.clear(); + let event = match xml_reader.read_event_into_async(&mut buf).await { + Ok(e) => e, + Err(e) => { + error!("Error reading epg XML event: {e}"); + break; + } + }; + + if let Some(flt) = &filter { + // Filter + match &event { + Event::Start(e) => { + if skip_depth.is_none() { + let should_skip = match e.name().as_ref() { + b"channel" => { + e.attributes() + .filter_map(Result::ok) + .find(|a| a.key.as_ref() == b"id") + .and_then(|a| a.unescape_value().ok()) + .is_some_and(|v| !flt.eq(v.as_ref())) + } + b"programme" => { + e.attributes() + .filter_map(Result::ok) + .find(|a| a.key.as_ref() == b"channel") + .and_then(|a| a.unescape_value().ok()) + .is_some_and(|v| !flt.eq(v.as_ref())) + } + _ => false, + }; + + if should_skip { + skip_depth = Some(1); + continue; + } + } else { + skip_depth = skip_depth.map(|d| d + 1); + continue; + } + } + Event::End(_) => { + if let Some(depth) = skip_depth { + if depth == 1 { + skip_depth = None; + } else { + skip_depth = Some(depth - 1); + } + continue; + } + } + Event::Empty(_) => { + if skip_depth.is_some() { + continue; + } + } + _ => {} + } + + if skip_depth.is_some() { + continue; + } + } + + match &event { + Event::Start(ref e) if offset_minutes != 0 && e.name().as_ref() == b"programme" => { // Modify the attributes let mut elem = BytesStart::new(EPG_TAG_PROGRAMME); for attr in e.attributes() { @@ -261,18 +328,18 @@ async fn serve_epg_with_rewrites( elem.push_attribute(attr); } Err(e) => { - error!("Error parsing attribute: {e}"); + error!("Error parsing epg attribute: {e}"); } } } // Write the modified start event if let Err(e) = xml_writer.write_event_async(Event::Start(elem)).await { - error!("Failed to write Start event: {e}"); + error!("Failed to write epg Start event: {e}"); break; } } - Ok(ref event @ (Event::Empty(ref e) | Event::Start(ref e))) if rewrite_urls && e.name().as_ref() == b"icon" => { + ref event @ (Event::Empty(ref e) | Event::Start(ref e)) if rewrite_urls && e.name().as_ref() == b"icon" => { // Modify the attributes let mut elem = BytesStart::new(EPG_TAG_ICON); for attr in e.attributes() { @@ -298,16 +365,11 @@ async fn serve_epg_with_rewrites( elem.push_attribute(attr); } Err(e) => { - error!("Error parsing attribute: {e}"); + error!("Error parsing epg attribute: {e}"); } } } - // Write the modified icon event - // if let Err(e) = xml_writer.write_event_async(Event::Start(elem)).await { - // error!("Failed to write Start event: {e}"); - // break; - // } let out_event = match event { Event::Empty(_) => Some(Event::Empty(elem)), Event::Start(_) => Some(Event::Start(elem)), @@ -315,28 +377,23 @@ async fn serve_epg_with_rewrites( }; if let Some(out) = out_event { if let Err(e) = xml_writer.write_event_async(out).await { - error!("Failed to write icon event: {e}"); + error!("Failed to write epg icon event: {e}"); break; } } } - Ok(Event::Decl(_) | Event::DocType(_)) => {}, - Ok(Event::Eof) => break, // End of file - Ok(event) => { + Event::Decl(_) | Event::DocType(_) => {}, + Event::Eof => break, // End of file + _ => { // Write any other event as is if let Err(e) = xml_writer.write_event_async(event).await { - error!("Failed to write event: {e}"); + error!("Failed to epg write event: {e}"); break; } } - Err(e) => { - error!("Error: {e}"); - break; - } } - - buf.clear(); } + buf.clear(); let mut encoder = xml_writer.into_inner(); if let Err(e) = encoder.shutdown().await { error!("Failed to shutdown epg gzip encoder: {e}"); @@ -386,7 +443,7 @@ async fn xmltv_api( return get_empty_epg_response(); }; - serve_epg(&app_state, &epg_path, &user, &target).await + serve_epg(&app_state, &epg_path, &user, &target, None).await } #[axum::debug_handler] diff --git a/backend/src/api/endpoints/xtream_api.rs b/backend/src/api/endpoints/xtream_api.rs index f45cefe4d..95044dba1 100644 --- a/backend/src/api/endpoints/xtream_api.rs +++ b/backend/src/api/endpoints/xtream_api.rs @@ -9,7 +9,7 @@ use crate::api::api_utils::{ }; use crate::api::api_utils::{redirect, try_result_not_found, try_option_bad_request, try_result_bad_request}; use crate::api::endpoints::hls_api::handle_hls_stream_request; -use crate::api::endpoints::xmltv_api::get_empty_epg_response; +use crate::api::endpoints::xmltv_api::{get_empty_epg_response, get_epg_path_for_target, serve_epg}; use crate::api::model::AppState; use crate::api::model::UserApiRequest; use crate::api::model::XtreamAuthorizationResponse; @@ -256,7 +256,7 @@ async fn xtream_player_api_stream( let input = try_option_bad_request!( app_state.app_config.get_input_by_name(pli.input_name.as_str()), true, - format!( "Cant find input for target {target_name}, context {}, stream_id {virtual_id}", stream_req.context) + format!( "Cant find input {} for target {target_name}, context {}, stream_id {virtual_id}", pli.input_name, stream_req.context) ); let (cluster, item_type) = if stream_req.context == ApiStreamContext::Timeshift { @@ -291,7 +291,7 @@ async fn xtream_player_api_stream( .into_response(); } - let stream_channel = create_stream_channel_with_type(&pli, item_type); + let stream_channel = create_stream_channel_with_type(target.id, &pli, item_type); if session.virtual_id == virtual_id && is_seek_request(cluster, req_headers).await { // partial request means we are in reverse proxy mode, seek happened @@ -380,7 +380,7 @@ async fn xtream_player_api_stream( .into_response(); } - let stream_channel = create_stream_channel_with_type(&pli, item_type); + let stream_channel = create_stream_channel_with_type(target.id, &pli, item_type); stream_response( fingerprint, @@ -432,8 +432,8 @@ async fn xtream_player_api_stream_with_token( .get_input_by_name(pli.input_name.as_str()), true, format!( - "Cant find input for target {target_name}, context {}, stream_id {}", - stream_req.context, pli.virtual_id + "Cant find input {} for target {target_name}, context {}, stream_id {}", + pli.input_name, stream_req.context, pli.virtual_id ) ); @@ -516,7 +516,7 @@ async fn xtream_player_api_stream_with_token( fingerprint, app_state, session_key.as_str(), - pli.to_stream_channel(), + pli.to_stream_channel(target.id), &stream_url, req_headers, &input, @@ -946,7 +946,7 @@ async fn xtream_player_api_timeshift_query_stream( async fn xtream_get_stream_info_response( app_state: &Arc, user: &ProxyUserCredentials, - target: &ConfigTarget, + target: &Arc, stream_id: &str, cluster: XtreamCluster, ) -> impl IntoResponse + Send { @@ -1029,7 +1029,7 @@ async fn xtream_get_stream_info_response( async fn xtream_get_short_epg( app_state: &Arc, user: &ProxyUserCredentials, - target: &ConfigTarget, + target: &Arc, stream_id: &str, limit: &str, ) -> impl IntoResponse + Send { @@ -1046,6 +1046,16 @@ async fn xtream_get_short_epg( target, None, ).await { + + let config = &app_state.app_config.config.load(); + if let Some(epg_path) = get_epg_path_for_target(config, target) { + if let Ok(exists) = tokio::fs::try_exists(&epg_path).await { + if exists { + return serve_epg(app_state, &epg_path, user, target, pli.epg_channel_id.clone()).await + } + } + } + if pli.provider_id > 0 { let input_name = &pli.input_name; if let Some(input) = app_state.app_config.get_input_by_name(input_name.as_str()) { @@ -1196,7 +1206,7 @@ async fn xtream_player_api_handle_content_action( async fn xtream_get_catchup_response( app_state: &Arc, - target: &ConfigTarget, + target: &Arc, stream_id: &str, start: &str, end: &str, diff --git a/backend/src/api/model/active_provider_manager.rs b/backend/src/api/model/active_provider_manager.rs index b81f03b2b..12c7bfc7b 100644 --- a/backend/src/api/model/active_provider_manager.rs +++ b/backend/src/api/model/active_provider_manager.rs @@ -2,41 +2,42 @@ use crate::api::model::provider_lineup_manager::{ProviderAllocation, ProviderLin use crate::api::model::{EventManager, ProviderConfig}; use crate::model::{AppConfig, ConfigInput}; use log::{debug, error}; +use crate::utils::{trace_if_enabled}; use shared::utils::{default_grace_period_millis, default_grace_period_timeout_secs}; use std::collections::{HashMap, HashSet}; use std::net::SocketAddr; use std::sync::Arc; use tokio::sync::RwLock; -pub type ProviderConnectionId = SocketAddr; +pub type ClientConnectionId = SocketAddr; #[derive(Debug, Clone)] pub struct ProviderHandle { - pub id: ProviderConnectionId, + pub client_id: ClientConnectionId, pub allocation: ProviderAllocation, } impl ProviderHandle { - pub fn new(id: ProviderConnectionId, allocation: ProviderAllocation) -> Self { - Self { id, allocation } + pub fn new(client_id: ClientConnectionId, allocation: ProviderAllocation) -> Self { + Self { client_id, allocation } } } #[derive(Debug, Clone)] struct SharedAllocation { allocation: ProviderAllocation, - connections: HashSet, + connections: HashSet, } #[derive(Debug, Clone, Default)] struct SharedConnections { by_key: HashMap, - key_by_addr: HashMap, + key_by_addr: HashMap, } #[derive(Debug, Clone, Default)] struct Connections { - single: HashMap, + single: HashMap, shared: SharedConnections, } @@ -74,9 +75,6 @@ impl ActiveProviderManager { } async fn acquire_connection_inner(&self, provider_or_input_name: &str, addr: &SocketAddr, force: bool) -> Option { - // Lock connections - let mut connections = self.connections.write().await; - // Call the specific acquisition function let allocation = if force { self.providers.force_exact_acquire_connection(provider_or_input_name).await @@ -88,12 +86,13 @@ impl ActiveProviderManager { ProviderAllocation::Exhausted => {} ProviderAllocation::Available(_) | ProviderAllocation::GracePeriod(_) => { let provider_name = allocation.get_provider_name().unwrap_or_default(); - + let mut connections = self.connections.write().await; if let Some(old) = connections.single.insert(*addr, allocation.clone()) { - crate::utils::trace_if_enabled!( + trace_if_enabled!( "register_connection: address {addr} already had a allocation for provider {:?} — forcing release on the old allocation", - old.get_provider_name().unwrap_or_default() - ); + old.get_provider_name().unwrap_or_default()); + + drop(connections); old.release().await; } @@ -128,37 +127,64 @@ impl ActiveProviderManager { } pub async fn release_connection(&self, addr: &SocketAddr) { - let mut connections = self.connections.write().await; + // Single connection + let single_allocation = { + let mut connections = self.connections.write().await; + connections.single.remove(addr) + }; - // try to release the single connection (not shared) - let handle = connections.single.remove(addr); - if let Some(allocation) = handle { - debug!("Released provider connection {:?} for {addr}", allocation.get_provider_name().unwrap_or_default()); + if let Some(allocation) = single_allocation { + debug!( + "Released provider connection {:?} for {addr}", + allocation.get_provider_name().unwrap_or_default() + ); allocation.release().await; return; } - let key = match connections.shared.key_by_addr.get(addr) { - Some(k) => k.clone(), - None => return, + // Shared connection + let shared_allocation = { + let mut connections = self.connections.write().await; + + let key = match connections.shared.key_by_addr.get(addr) { + Some(k) => k.clone(), + None => return, // no shared connection + }; + + // 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 + } }; - let mut released = false; - if let Some(connections) = connections.shared.by_key.get_mut(&key) { - if connections.connections.remove(addr) && connections.connections.is_empty() { - connections.allocation.release().await; - released = true; - } - } - - if released { - connections.shared.key_by_addr.remove(addr); - connections.shared.by_key.remove(&key); + // release allocation + if let Some(allocation) = shared_allocation { + 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.id).await; + self.release_connection(&handle.client_id).await; } pub async fn make_shared_connection(&self, addr: &SocketAddr, key: &str) { @@ -178,7 +204,7 @@ impl ActiveProviderManager { shared_allocation.connections.insert(*addr); connections.shared.key_by_addr.insert(*addr, key.to_string()); } else { - error!("Failed to add shared connection for {addr}: url: {key:?} not found"); + error!("Failed to add shared connection for {addr}: url: {key:?} not found"); } } diff --git a/backend/src/api/model/connection_manager.rs b/backend/src/api/model/connection_manager.rs index f1e692546..4dab2e5af 100644 --- a/backend/src/api/model/connection_manager.rs +++ b/backend/src/api/model/connection_manager.rs @@ -58,7 +58,7 @@ impl ConnectionManager { pub async fn release_provider_handle(&self, provider_handle: Option) { if let Some(handle) = provider_handle { - self.release_provider_connection(&handle.id).await; + self.release_provider_connection(&handle.client_id).await; } } diff --git a/backend/src/api/model/provider_config.rs b/backend/src/api/model/provider_config.rs index 09e604626..e69376aa4 100644 --- a/backend/src/api/model/provider_config.rs +++ b/backend/src/api/model/provider_config.rs @@ -1,5 +1,5 @@ use std::fmt; -use crate::model::{ConfigInput, ConfigInputAlias, InputUserInfo}; +use crate::model::{is_input_expired, ConfigInput, ConfigInputAlias, InputUserInfo}; use jsonwebtoken::get_current_timestamp; use log::{debug}; use std::ops::Deref; @@ -42,6 +42,7 @@ pub struct ProviderConfig { pub input_type: InputType, max_connections: usize, priority: i16, + exp_date: Option, connection: RwLock, on_connection_change: ProviderConnectionChangeCallback, } @@ -57,7 +58,8 @@ impl fmt::Display for ProviderConfig { write!(f, ", priority: {}", self.priority)?; write_if_some!(f, self, ", username: " => username, - ", password: " => password + ", password: " => password, + ", exp_date: " => exp_date ); write!(f, "}}")?; Ok(()) @@ -80,6 +82,7 @@ impl PartialEq for ProviderConfig { && self.input_type == other.input_type && self.max_connections == other.max_connections && self.priority == other.priority + && self.exp_date == other.exp_date // Note: self.connection is skipped } } @@ -90,7 +93,7 @@ macro_rules! modify_connections { $self.notify_connection_change($guard.current_connections); }}; ($self:ident, $guard:ident, -1) => {{ - $guard.current_connections -= 1; + $guard.current_connections = $guard.current_connections.saturating_sub(1); $self.notify_connection_change($guard.current_connections); }}; } @@ -109,6 +112,7 @@ impl ProviderConfig { input_type: cfg.input_type, max_connections: cfg.max_connections as usize, priority: cfg.priority, + exp_date: cfg.exp_date, connection: RwLock::new(get_connection.and_then(|f| f(cfg.name.as_str())).map_or_else(Default::default, Clone::clone)), on_connection_change } @@ -127,6 +131,7 @@ impl ProviderConfig { input_type: cfg.input_type, max_connections: alias.max_connections as usize, priority: alias.priority, + exp_date: alias.exp_date, connection: RwLock::new(get_connection.and_then(|f| f(alias.name.as_str())).map_or_else(Default::default, Clone::clone)), on_connection_change, } @@ -178,12 +183,20 @@ impl ProviderConfig { // !self.is_exhausted() // } - async fn force_allocate(&self) { + async fn force_allocate(&self) -> bool { + if is_input_expired(self.exp_date) { + return false; + } let mut guard = self.connection.write().await; modify_connections!(self, guard, +1); + true } async fn try_allocate(&self, grace: bool, grace_period_timeout_secs: u64) -> ProviderConfigAllocation { + if is_input_expired(self.exp_date) { + return ProviderConfigAllocation::Exhausted; + } + let mut guard = self.connection.write().await; if self.max_connections == 0 { modify_connections!(self, guard, +1); @@ -220,6 +233,10 @@ impl ProviderConfig { // is intended to use with redirects, to cycle through provider // do not increment and connection counter! async fn get_next(&self, grace: bool, grace_period_timeout_secs: u64) -> bool { + if is_input_expired(self.exp_date) { + return false; + } + if self.max_connections == 0 { return true; } @@ -288,8 +305,11 @@ impl ProviderConfigWrapper { } pub async fn force_allocate(&self) -> ProviderAllocation { - self.inner.force_allocate().await; - ProviderAllocation::new_available(Arc::clone(&self.inner)) + if self.inner.force_allocate().await { + ProviderAllocation::new_available(Arc::clone(&self.inner)) + } else { + ProviderAllocation::Exhausted + } } pub async fn try_allocate(&self, grace: bool, grace_period_timeout_secs: u64) -> ProviderAllocation { diff --git a/backend/src/api/model/provider_lineup_manager.rs b/backend/src/api/model/provider_lineup_manager.rs index 88834306a..9ccd4ae08 100644 --- a/backend/src/api/model/provider_lineup_manager.rs +++ b/backend/src/api/model/provider_lineup_manager.rs @@ -555,6 +555,7 @@ impl ProviderLineupManager { || a.username != b.username || a.password != b.password || a.url != b.url + || a.exp_date != b.exp_date { return true; } @@ -577,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/api/model/streams/active_client_stream.rs b/backend/src/api/model/streams/active_client_stream.rs index 2819d9f99..fea261151 100644 --- a/backend/src/api/model/streams/active_client_stream.rs +++ b/backend/src/api/model/streams/active_client_stream.rs @@ -1,4 +1,4 @@ -use crate::api::model::{AppState, CustomVideoStreamType, ProviderHandle, StreamDetails}; +use crate::api::model::{AppState, ConnectionManager, CustomVideoStreamType, ProviderHandle, StreamDetails}; use crate::api::model::BoxedProviderStream; use crate::api::model::StreamError; use crate::api::model::TimedClientStream; @@ -30,6 +30,7 @@ pub(in crate::api) struct ActiveClientStream { provider_handle: Option, custom_video: (Option, Option), waker: Option>, + connection_manager: Arc, } impl ActiveClientStream { @@ -95,6 +96,7 @@ impl ActiveClientStream { send_custom_stream_flag: grace_stop_flag, custom_video, waker, + connection_manager: Arc::clone(&app_state.connection_manager) } } @@ -218,3 +220,13 @@ impl Stream for ActiveClientStream { Poll::Ready(None) } } + +impl Drop for ActiveClientStream { + fn drop(&mut self) { + let mgr = Arc::clone(&self.connection_manager); + let hndl = self.provider_handle.take(); + tokio::spawn(async move { + mgr.release_provider_handle(hndl).await; + }); + } +} \ No newline at end of file diff --git a/backend/src/api/model/streams/provider_stream_factory.rs b/backend/src/api/model/streams/provider_stream_factory.rs index af508e5f3..841c85727 100644 --- a/backend/src/api/model/streams/provider_stream_factory.rs +++ b/backend/src/api/model/streams/provider_stream_factory.rs @@ -354,11 +354,8 @@ async fn handle_channel_unavailable_stream(app_state: &Arc, app_state.connection_manager.release_provider_connection(&stream_options.addr).await; if let (Some(boxed_provider_stream), response_info) = - create_channel_unavailable_stream( - &app_state.app_config, - &get_response_headers(stream_options.get_headers()), - StatusCode::SERVICE_UNAVAILABLE, - ) + create_channel_unavailable_stream(&app_state.app_config,&get_response_headers(stream_options.get_headers()), + StatusCode::SERVICE_UNAVAILABLE) { Ok(Some((boxed_provider_stream, response_info))) } else { @@ -383,37 +380,23 @@ async fn get_provider_stream( } Ok(None) => { if connect_err > ERR_MAX_RETRY_COUNT { - warn!( - "The stream could be unavailable. {}", - sanitize_sensitive_info(stream_options.get_url().as_str()) - ); + warn!("The stream could be unavailable. {}", sanitize_sensitive_info(stream_options.get_url().as_str())); + break; } } Err(status) => { debug!("Provider stream response error status response : {status}"); - if status == StatusCode::FORBIDDEN - || status == StatusCode::SERVICE_UNAVAILABLE - || status == StatusCode::UNAUTHORIZED - { - warn!( - "The stream could be unavailable. ({status}) {}", - sanitize_sensitive_info(stream_options.get_url().as_str()) - ); - stream_options.cancel_reconnect(); - return Err(status); + if matches!(status, StatusCode::FORBIDDEN | StatusCode::SERVICE_UNAVAILABLE | StatusCode::UNAUTHORIZED) { + warn!("The stream could be unavailable. ({status}) {}",sanitize_sensitive_info(stream_options.get_url().as_str())); + break; } if connect_err > ERR_MAX_RETRY_COUNT { - warn!( - "The stream could be unavailable. ({status}) {}", - sanitize_sensitive_info(stream_options.get_url().as_str()) - ); + warn!("The stream could be unavailable. ({status}) {}",sanitize_sensitive_info(stream_options.get_url().as_str())); + break; } } } - if !stream_options.should_continue() { - return Err(StatusCode::SERVICE_UNAVAILABLE); - } - if connect_err > ERR_MAX_RETRY_COUNT { + if !stream_options.should_continue() || connect_err > ERR_MAX_RETRY_COUNT { break; } if start.elapsed().as_secs() > RETRY_SECONDS { @@ -425,16 +408,11 @@ async fn get_provider_stream( } connect_err += 1; tokio::time::sleep(Duration::from_millis(50)).await; - debug_if_enabled!( - "Reconnecting stream {}", - sanitize_sensitive_info(url.as_str()) - ); + debug_if_enabled!("Reconnecting stream {}", sanitize_sensitive_info(url.as_str())); } - debug_if_enabled!( - "Stopped reconnecting stream {}", - sanitize_sensitive_info(url.as_str()) - ); + debug_if_enabled!("Stopped reconnecting stream {}", sanitize_sensitive_info(url.as_str())); stream_options.cancel_reconnect(); + app_state.connection_manager.release_provider_connection(&stream_options.addr).await; Err(StatusCode::SERVICE_UNAVAILABLE) } @@ -493,7 +471,20 @@ pub async fn create_provider_stream( if continue_streaming.is_active() { match get_provider_stream(&app_state_clone, &client, &stream_opts).await { Ok(Some((stream, _info))) => Some((stream, ())), - Ok(None) => None, + Ok(None) => { + app_state_clone.connection_manager.release_provider_connection(&stream_opts.addr).await; + continue_streaming.notify(); + if let (Some(boxed_provider_stream), _response_info) = + create_channel_unavailable_stream( + &app_state_clone.app_config, + &get_response_headers(stream_opts.get_headers()), + StatusCode::SERVICE_UNAVAILABLE, + ) + { + return Some((boxed_provider_stream, ())); + } + None + } Err(status) => { app_state_clone.connection_manager.release_provider_connection(&stream_opts.addr).await; continue_streaming.notify(); @@ -510,6 +501,7 @@ pub async fn create_provider_stream( } } } else { + app_state_clone.connection_manager.release_provider_connection(&stream_opts.addr).await; None } } diff --git a/backend/src/api/model/streams/shared_stream_manager.rs b/backend/src/api/model/streams/shared_stream_manager.rs index 5c9d31faa..2d432ecea 100644 --- a/backend/src/api/model/streams/shared_stream_manager.rs +++ b/backend/src/api/model/streams/shared_stream_manager.rs @@ -20,6 +20,7 @@ use std::pin::Pin; use std::task::{Context, Poll}; use tokio::sync::mpsc::Sender; use tokio::sync::{mpsc, Mutex, RwLock}; +use tokio::time::{sleep, Duration, Instant}; use tokio_stream::wrappers::ReceiverStream; use tokio_util::sync::CancellationToken; @@ -96,7 +97,7 @@ impl BurstBuffer { } pub fn push(&mut self, packet: Arc) { - while self.current_bytes > self.buffer_size { + while self.current_bytes + packet.len() > self.buffer_size { if let Some(popped) = self.buffer.pop_front() { self.current_bytes -= popped.len(); } else { @@ -135,6 +136,7 @@ pub struct SharedStreamState { broadcaster: tokio::sync::broadcast::Sender, stop_token: CancellationToken, burst_buffer: Arc>, + task_handles: Mutex>>, } impl SharedStreamState { @@ -151,6 +153,7 @@ impl SharedStreamState { broadcaster, stop_token: CancellationToken::new(), burst_buffer: Arc::new(Mutex::new(BurstBuffer::new(burst_buffer_size_in_bytes))), + task_handles: Mutex::new(Vec::new()), } } @@ -158,6 +161,7 @@ impl SharedStreamState { let (client_tx, client_rx) = mpsc::channel(self.buf_size); let mut broadcast_rx = self.broadcaster.subscribe(); let cancel_token = CancellationToken::new(); + { let mut subs = self.subscribers.write().await; subs.insert(*addr, cancel_token.clone()); @@ -169,8 +173,13 @@ impl SharedStreamState { let burst_buffer_for_log = Arc::clone(&self.burst_buffer); let yield_counter = YIELD_COUNTER; + // If a client stops streaming (for example presses + let timeout_duration = Duration::from_secs(300); // 5 minutes + let mut last_active = Instant::now(); + let address = *addr; - tokio::spawn(async move { + let handle = tokio::spawn(async move { + // initial burst buffer let snapshot = { let buffer = burst_buffer.lock().await; buffer.snapshot() @@ -180,47 +189,63 @@ impl SharedStreamState { let mut loop_cnt = 0; loop { tokio::select! { - biased; + biased; - () = cancel_token.cancelled() => { - debug!("Client disconnected from shared stream: {address}"); + // canceled + () = cancel_token.cancelled() => { + debug!("Client disconnected from shared stream: {address}"); + break; + } + + // timeout handling + () = sleep(Duration::from_secs(1)) => { + if last_active.elapsed() > timeout_duration { + debug!("Client timed out due to inactivity: {address}"); + cancel_token.cancel(); break; } - result = broadcast_rx.recv() => { - match result { - Ok(data) => { - if let Err(err) = client_tx.send(data).await { - debug!("Shared stream client send error: {address} {err}"); - break; - } - loop_cnt += 1; - if loop_cnt >= yield_counter { - tokio::task::yield_now().await; - loop_cnt = 0; - } + } + + // receive broadcast data + result = broadcast_rx.recv() => { + match result { + Ok(data) => { + // Wenn der Client pausiert, einfach skippen oder warten + if client_tx_clone.is_closed() { + continue; } - Err(tokio::sync::broadcast::error::RecvError::Lagged(skipped)) => { - let buffered_bytes = { - let buffer = burst_buffer_for_log.lock().await; - buffer.current_bytes - }; - warn!("Shared stream client lagged behind {address}. Skipped {skipped} messages (buffered {buffered_bytes} bytes, yield counter {yield_counter})"); - loop_cnt += 1; - if loop_cnt >= yield_counter { - tokio::task::yield_now().await; - loop_cnt = 0; - } + + if let Err(err) = client_tx.send(data).await { + debug!("Shared stream client send error: {address} {err}"); + break; + } + loop_cnt += 1; + last_active = Instant::now(); + + if loop_cnt >= yield_counter { + tokio::task::yield_now().await; + loop_cnt = 0; } - Err(_) => break, } + Err(tokio::sync::broadcast::error::RecvError::Lagged(skipped)) => { + let buffered_bytes = { + let buffer = burst_buffer_for_log.lock().await; + buffer.current_bytes + }; + warn!("Shared stream client lagged behind {address}. Skipped {skipped} messages (buffered {buffered_bytes} bytes, yield counter {yield_counter})"); + } + Err(_) => break, } } } + } + manager.release_connection(&address, false).await; }); - let provider = self.provider_guard.as_ref().and_then(|h| h.allocation.get_provider_name()); + self.task_handles.lock().await.push(handle); + let provider = self.provider_guard.as_ref().and_then(|h| h.allocation.get_provider_name()); (convert_stream(ReceiverStream::new(client_rx).boxed()), provider) } @@ -325,34 +350,37 @@ impl SharedStreamManager { self.get_shared_state(stream_url).await.map(|s| s.headers.clone()) } - async fn unregister(&self, stream_url: &str, send_stop_signal: bool) - { - let shared_state = { + async fn unregister(&self, stream_url: &str, send_stop_signal: bool) { + let shared_state_opt = { let mut shared_streams = self.shared_streams.write().await; - let shared_state = shared_streams.by_key.remove(stream_url); - { - let remove_keys: Vec = shared_streams.key_by_addr.iter() - .filter_map(|(addr, url)| if url == stream_url { Some(*addr) } else { None }) - .collect(); - for k in remove_keys { - shared_streams.key_by_addr.remove(&k); - } + + let remove_keys: Vec = shared_streams.key_by_addr + .iter() + .filter_map(|(addr, url)| if url == stream_url { Some(*addr) } else { None }) + .collect(); + for k in remove_keys { + shared_streams.key_by_addr.remove(&k); } - shared_state + + shared_streams.by_key.remove(stream_url) }; - if let Some(shared_state) = shared_state { + if let Some(shared_state) = shared_state_opt { let remaining = shared_state.subscribers.read().await.len(); debug_if_enabled!("Unregistering shared stream {} (remaining_subscribers={remaining}, send_stop_signal={send_stop_signal})", - sanitize_sensitive_info(stream_url)); + sanitize_sensitive_info(stream_url)); + + for handle in shared_state.task_handles.lock().await.drain(..) { + handle.abort(); + } if let Some(provider_handle) = &shared_state.provider_guard { self.provider_manager.release_handle(provider_handle).await; } - if send_stop_signal { + if send_stop_signal || remaining == 0 { trace_if_enabled!("Sending shared stream stop signal {}", sanitize_sensitive_info(stream_url)); - let () = shared_state.stop_token.cancel(); + shared_state.stop_token.cancel(); } } } @@ -360,54 +388,65 @@ impl SharedStreamManager { pub async fn release_connection(&self, addr: &SocketAddr, send_stop_signal: bool) { let (stream_url, shared_state) = { let shared_streams = self.shared_streams.read().await; - if let Some(stream_url) = shared_streams.key_by_addr.get(addr) { - (Some(stream_url.clone()), shared_streams.by_key.get(stream_url).cloned()) - } else { - (None, None) - } - }; + if let Some(stream_url) = shared_streams.key_by_addr.get(addr) { + (Some(stream_url.clone()), shared_streams.by_key.get(stream_url).cloned()) + } else { + (None, None) + } + }; if let Some(state) = shared_state { let (tx, is_empty, remaining) = { let mut subs = state.subscribers.write().await; let tx = subs.remove(addr); let is_empty = subs.is_empty(); - ( - if send_stop_signal { tx } else { None }, - is_empty, - subs.len(), - ) + (tx, is_empty, subs.len()) }; debug!("Shared stream subscriber removed {addr}; remaining subscribers={remaining}"); if is_empty { if let Some(url) = stream_url.as_ref() { - debug_if_enabled!("No subscribers remain for {} after removing {addr}", sanitize_sensitive_info(url) - ); + debug_if_enabled!( + "No subscribers remain for {} after removing {addr}", + sanitize_sensitive_info(url) + ); self.unregister(url, send_stop_signal).await; } } if let Some(client_stop_signal) = tx { - let () = client_stop_signal.cancel(); + client_stop_signal.cancel(); } } } - async fn subscribe_stream(&self, stream_url: &str, addr: &SocketAddr, manager: Arc) -> Option<(BoxedProviderStream, Option)> { - let mut shared_streams = self.shared_streams.write().await; - let shared_state_opt = shared_streams.by_key.get(stream_url).cloned(); - match shared_state_opt { - Some(shared_state) => { - debug_if_enabled!("Responding to existing shared client stream {addr} {}", sanitize_sensitive_info(stream_url)); + async fn subscribe_stream( + &self, + stream_url: &str, + addr: &SocketAddr, + manager: Arc, + ) -> Option<(BoxedProviderStream, Option)> { + let shared_state_opt = { + let shared_streams = self.shared_streams.read().await; + shared_streams.by_key.get(stream_url).cloned() + }; + + if let Some(shared_state) = shared_state_opt { + { + let mut shared_streams = self.shared_streams.write().await; shared_streams.key_by_addr.insert(*addr, stream_url.to_owned()); - Some(shared_state.subscribe(addr, manager).await) } - None => None, + + debug_if_enabled!("Responding to existing shared client stream {addr} {}",sanitize_sensitive_info(stream_url) + ); + Some(shared_state.subscribe(addr, manager).await) + } else { + None } } + async fn register(&self, addr: &SocketAddr, stream_url: &str, shared_state: Arc) { let mut shared_streams = self.shared_streams.write().await; shared_streams.by_key.insert(stream_url.to_string(), shared_state); diff --git a/backend/src/model/config/input.rs b/backend/src/model/config/input.rs index 4524b66c9..e2d8a1040 100644 --- a/backend/src/model/config/input.rs +++ b/backend/src/model/config/input.rs @@ -1,14 +1,16 @@ use crate::model::{macros, EpgConfig}; -use shared::error::{TuliproxError}; -use shared::{check_input_connections, info_err, write_if_some}; +use crate::utils::get_csv_file_path; +use chrono::Utc; +use log::warn; +use shared::error::TuliproxError; use shared::model::{ConfigInputAliasDto, ConfigInputDto, ConfigInputOptionsDto, InputFetchMethod, InputType, StagedInputDto}; use shared::utils::{get_base_url_from_str, get_credentials_from_url}; -use shared::{check_input_credentials}; +use shared::{check_input_connections, info_err, write_if_some}; +use shared::check_input_credentials; use std::collections::HashMap; use std::fmt; use std::path::PathBuf; use url::Url; -use crate::utils::{get_csv_file_path}; #[allow(clippy::struct_excessive_bools)] #[derive(Debug, Clone)] @@ -101,6 +103,7 @@ pub struct ConfigInputAlias { pub password: Option, pub priority: i16, pub max_connections: u16, + pub exp_date: Option, } macros::from_impl!(ConfigInputAlias); @@ -114,6 +117,7 @@ impl From<&ConfigInputAliasDto> for ConfigInputAlias { password: dto.password.clone(), priority: dto.priority, max_connections: dto.max_connections, + exp_date: dto.exp_date, } } } @@ -136,6 +140,7 @@ pub struct ConfigInput { pub max_connections: u16, pub method: InputFetchMethod, pub staged: Option, + pub exp_date: Option, pub t_batch_url: Option, } @@ -150,6 +155,12 @@ impl ConfigInput { return Err(info_err!("Staged input can only be from type m3u or xtream".to_owned())); } } + + if is_input_expired(self.exp_date) { + warn!("Account {} expired for provider: {}", self.username.as_ref().map_or("?", |s| s.as_str()), self.name); + self.enabled = false; + } + Ok(batch_file_path) } @@ -180,10 +191,16 @@ impl ConfigInput { InputType::Xtream }; - self.t_batch_url= Some(self.url.clone()); + self.t_batch_url = Some(self.url.clone()); let file_path = get_csv_file_path(self.url.as_str()).ok(); if let Some(aliases) = self.aliases.as_mut() { + for alias in aliases.iter() { + if is_input_expired(alias.exp_date) { + warn!("Alias-Account {} expired for provider: {}", alias.username.as_ref().map_or("?", |s| s.as_str()), alias.name); + } + } + if !aliases.is_empty() { let mut first = aliases.remove(0); self.id = first.id; @@ -223,6 +240,7 @@ impl ConfigInput { max_connections: alias.max_connections, method: self.method, staged: None, + exp_date: None, t_batch_url: None, } } @@ -247,8 +265,9 @@ impl From<&ConfigInputDto> for ConfigInput { priority: dto.priority, max_connections: dto.max_connections, method: dto.method, - t_batch_url: None, + exp_date: dto.exp_date, staged: dto.staged.as_ref().map(StagedInput::from), + t_batch_url: None, } } } @@ -277,3 +296,13 @@ impl fmt::Display for ConfigInput { Ok(()) } } + +pub fn is_input_expired(exp_date: Option) -> bool { + match exp_date { + Some(ts) => { + let now = Utc::now().timestamp(); + ts <= now + } + None => false, + } +} \ No newline at end of file diff --git a/backend/src/model/xmltv.rs b/backend/src/model/xmltv.rs index 17ec99793..d4a5f131c 100644 --- a/backend/src/model/xmltv.rs +++ b/backend/src/model/xmltv.rs @@ -230,6 +230,7 @@ pub fn get_attr_value(attr: &quick_xml::events::attributes::Attribute) -> Option attr.unescape_value().ok().map(|v| v.to_string()) } +// This function filters a timeslot starting from yesterday. #[allow(clippy::too_many_lines)] async fn parse_xmltv_for_web_ui(reader: R) -> Result { @@ -246,11 +247,11 @@ async fn parse_xmltv_for_web_ui(reader: R) -> Resul // only 1 day old epg let now = Utc::now(); - let yesterday_start = Utc.with_ymd_and_hms(now.year(), now.month(), now.day(), 0, 0, 0).unwrap() + let yesterday_start = Utc.with_ymd_and_hms(now.year(), now.month(), now.day(), 0, 0, 0) + .single().expect("Current date at midnight should always be valid") - chrono::Duration::days(1); let threshold_ts = yesterday_start.timestamp(); - loop { match reader.read_event_into_async(&mut buf).await { Ok(Event::Empty(e) | Event::Start(e)) => { diff --git a/backend/src/model/xtream.rs b/backend/src/model/xtream.rs index 550ed2638..f740b2eb9 100644 --- a/backend/src/model/xtream.rs +++ b/backend/src/model/xtream.rs @@ -2,12 +2,18 @@ use crate::model::{AppConfig, ProxyUserCredentials}; use crate::model::{ConfigTarget, XtreamTargetOutput}; use serde::{Deserialize, Deserializer, Serialize}; use serde_json::{Map, Value}; -use shared::model::{xtream_const, PlaylistItem, XtreamPlaylistItem}; +use shared::model::{xtream_const, PlaylistItem, ProxyUserStatus, XtreamPlaylistItem}; use shared::model::{ClusterFlags, PlaylistEntry, XtreamCluster}; use shared::utils::{deserialize_as_option_string, deserialize_as_string, deserialize_as_string_array, deserialize_number_from_string, 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, +} + #[derive(Deserialize, Default)] pub struct XtreamCategory { #[serde(deserialize_with = "deserialize_as_string")] diff --git a/backend/src/repository/epg_repository.rs b/backend/src/repository/epg_repository.rs index 9faad0e5e..a1872015b 100644 --- a/backend/src/repository/epg_repository.rs +++ b/backend/src/repository/epg_repository.rs @@ -7,20 +7,30 @@ use shared::error::{notify_err, TuliproxError}; use std::path::Path; use tokio::io::AsyncWriteExt; + +// Due to an error in quick_xml we cant write doc type through event. The quotes are escaped and the xml file is invalid. +// +// // XML Header +// writer.write_event_async(quick_xml::events::Event::Decl(quick_xml::events::BytesDecl::new("1.0", Some("utf-8"), None))) +// .await.map_err(|e| notify_err!(format!("failed to write XML header: {}", e)))?; +// +// // DOCTYPE +// writer.write_event_async(quick_xml::events::Event::DocType(quick_xml::events::BytesText::new(r#"tv SYSTEM "xmltv.dtd""#))) +// .await.map_err(|e| notify_err!(format!("failed to write doctype: {}", e)))?; pub async fn epg_write_file(target: &ConfigTarget, epg: &Epg, path: &Path) -> Result<(), TuliproxError> { let file = tokio::fs::File::create(path).await .map_err(|e| notify_err!(format!("failed to create epg file: {}", e)))?; - let buf_writer = tokio::io::BufWriter::new(file); + let mut buf_writer = tokio::io::BufWriter::new(file); + + // Work-Around BytesText DocType escape, see below + buf_writer.write_all(b"\n").await + .map_err(|e| notify_err!(format!("failed to write XML header: {}", e)))?; + + buf_writer.write_all(b"\n").await + .map_err(|e| notify_err!(format!("failed to write doctype: {}", e)))?; + let mut writer = quick_xml::writer::Writer::new(buf_writer); - // XML Header - writer.write_event_async(quick_xml::events::Event::Decl(quick_xml::events::BytesDecl::new("1.0", Some("utf-8"), None))) - .await.map_err(|e| notify_err!(format!("failed to write XML header: {}", e)))?; - - // DOCTYPE - writer.write_event_async(quick_xml::events::Event::DocType(quick_xml::events::BytesText::new(r#"tv SYSTEM "xmltv.dtd""#))) - .await.map_err(|e| notify_err!(format!("failed to write doctype: {}", e)))?; - // EPG Content epg.write_to_async(&mut writer).await.map_err(|e| notify_err!(format!("failed to write epg: {}", e)))?; diff --git a/backend/src/repository/user_repository.rs b/backend/src/repository/user_repository.rs index d1c58b1c7..8efd43f03 100644 --- a/backend/src/repository/user_repository.rs +++ b/backend/src/repository/user_repository.rs @@ -8,7 +8,7 @@ use crate::utils; use crate::utils::json_write_documents_to_file; use chrono::Local; use log::error; -use shared::model::{PlaylistBouquetDto, ProxyType, ProxyUserStatus, TargetBouquetDto, TargetType, XtreamCluster}; +use shared::model::{PlaylistBouquetDto, PlaylistClusterBouquetDto, ProxyType, ProxyUserStatus, TargetType, XtreamCluster}; use std::collections::{HashMap, HashSet}; use std::io::Error; use std::path::{Path, PathBuf}; @@ -256,7 +256,6 @@ async fn save_xtream_user_bouquet_for_target(config: &Config, target_name: &str, XtreamCluster::Series => user_get_series_bouquet_path(storage_path, TargetType::Xtream), }; - if let Some(bouquet_categories) = bouquet { if let Some(xtream_categories) = xtream_get_playlist_categories(config, target_name, cluster).await { let filtered: Vec = xtream_categories.iter().filter(|p| bouquet_categories.contains(&p.name)).cloned().collect(); @@ -294,27 +293,23 @@ async fn save_m3u_user_bouquet_for_target(storage_path: &Path, target: TargetTyp Ok(()) } -async fn save_user_bouquet_for_target(config: &Config, target_name: &str, storage_path: &Path, target: TargetType, bouquet: &TargetBouquetDto) -> Result<(), Error> { +async fn save_user_bouquet_for_target(config: &Config, target_name: &str, storage_path: &Path, target: TargetType, bouquet: Option<&PlaylistClusterBouquetDto>) -> Result<(), Error> { if target == TargetType::Xtream { - save_xtream_user_bouquet_for_target(config, target_name, storage_path, XtreamCluster::Live, bouquet.live.as_ref()).await?; - save_xtream_user_bouquet_for_target(config, target_name, storage_path, XtreamCluster::Video, bouquet.vod.as_ref()).await?; - save_xtream_user_bouquet_for_target(config, target_name, storage_path, XtreamCluster::Series, bouquet.series.as_ref()).await?; + save_xtream_user_bouquet_for_target(config, target_name, storage_path, XtreamCluster::Live, bouquet.and_then(|b| b.live.as_ref())).await?; + save_xtream_user_bouquet_for_target(config, target_name, storage_path, XtreamCluster::Video, bouquet.and_then(|b| b.vod.as_ref())).await?; + save_xtream_user_bouquet_for_target(config, target_name, storage_path, XtreamCluster::Series, bouquet.and_then(|b| b.series.as_ref())).await?; } else { - save_m3u_user_bouquet_for_target(storage_path, target, XtreamCluster::Live, bouquet.live.as_ref()).await?; - save_m3u_user_bouquet_for_target(storage_path, target, XtreamCluster::Video, bouquet.vod.as_ref()).await?; - save_m3u_user_bouquet_for_target(storage_path, target, XtreamCluster::Series, bouquet.series.as_ref()).await?; + save_m3u_user_bouquet_for_target(storage_path, target, XtreamCluster::Live, bouquet.and_then(|b| b.live.as_ref())).await?; + save_m3u_user_bouquet_for_target(storage_path, target, XtreamCluster::Video, bouquet.and_then(|b| b.vod.as_ref())).await?; + save_m3u_user_bouquet_for_target(storage_path, target, XtreamCluster::Series, bouquet.and_then(|b| b.series.as_ref())).await?; } Ok(()) } pub async fn save_user_bouquet(cfg: &Config, target_name: &str, username: &str, bouquet: &PlaylistBouquetDto) -> Result<(), Error> { if let Some(storage_path) = ensure_user_storage_path(cfg, username) { - if let Some(xb) = &bouquet.xtream { - save_user_bouquet_for_target(cfg, target_name, &storage_path, TargetType::Xtream, xb).await?; - } - if let Some(mb) = &bouquet.m3u { - save_user_bouquet_for_target(cfg, target_name, &storage_path, TargetType::M3u, mb).await?; - } + save_user_bouquet_for_target(cfg, target_name, &storage_path, TargetType::Xtream, bouquet.xtream.as_ref()).await?; + save_user_bouquet_for_target(cfg, target_name, &storage_path, TargetType::M3u, bouquet.m3u.as_ref()).await?; Ok(()) } else { Err(Error::new(std::io::ErrorKind::NotFound, format!("User config path not found for user {username}"))) @@ -391,7 +386,6 @@ pub async fn user_get_bouquet_filter(config: &Config, username: &str, category_i XtreamCluster::Series => user_get_series_bouquet(config, username, target).await, }; - match bouquet { None => None, Some(bouquet_categories) => { diff --git a/backend/src/utils/file/csv_input_reader.rs b/backend/src/utils/file/csv_input_reader.rs index 776a41156..b64e52f12 100644 --- a/backend/src/utils/file/csv_input_reader.rs +++ b/backend/src/utils/file/csv_input_reader.rs @@ -7,7 +7,7 @@ use std::io::{BufRead, Cursor, Error}; use std::path::PathBuf; use url::Url; use shared::model::{ConfigInputAliasDto, InputType}; -use shared::utils::{get_credentials_from_url, trim_last_slash}; +use shared::utils::{get_credentials_from_url, parse_timestamp, trim_last_slash}; use crate::utils::request::get_local_file_content; const CSV_SEPARATOR: char = ';'; @@ -18,8 +18,9 @@ const FIELD_URL: &str = "url"; const FIELD_NAME: &str = "name"; const FIELD_USERNAME: &str = "username"; const FIELD_PASSWORD: &str = "password"; +const FIELD_EXP_DATE: &str = "exp_date"; const FIELD_UNKNOWN: &str = "?"; -const DEFAULT_COLUMNS: &[&str] = &[FIELD_URL, FIELD_MAX_CON, FIELD_PRIO, FIELD_NAME, FIELD_USERNAME, FIELD_PASSWORD]; +const DEFAULT_COLUMNS: &[&str] = &[FIELD_URL, FIELD_MAX_CON, FIELD_PRIO, FIELD_NAME, FIELD_USERNAME, FIELD_PASSWORD, FIELD_EXP_DATE]; fn csv_assign_mandatory_fields(alias: &mut ConfigInputAliasDto, input_type: InputType) { if !alias.url.is_empty() { @@ -86,6 +87,12 @@ fn csv_assign_config_input_column(config_input: &mut ConfigInputAliasDto, header FIELD_PASSWORD => { config_input.password = Some(value.to_string()); } + FIELD_EXP_DATE => { + config_input.exp_date = parse_timestamp(value).unwrap_or_else(|e| { + error!("Failed to parse exp_date '{value}': {e}"); + None + }); + } _ => {} } } @@ -117,6 +124,7 @@ pub fn csv_read_inputs_from_reader(batch_input_type: InputType, reader: impl Buf FIELD_NAME => FIELD_NAME, FIELD_USERNAME => FIELD_USERNAME, FIELD_PASSWORD => FIELD_PASSWORD, + FIELD_EXP_DATE => FIELD_EXP_DATE, _ => { error!("Field {s} is unsupported for csv input"); FIELD_UNKNOWN @@ -135,6 +143,7 @@ pub fn csv_read_inputs_from_reader(batch_input_type: InputType, reader: impl Buf password: None, priority: 0, max_connections: 1, + exp_date: None, }; let columns: Vec<&str> = line.split(CSV_SEPARATOR).collect(); @@ -193,9 +202,9 @@ http://hd.providerline.com/get.php?username=user4&password=user4&type=m3u_plus;i "; const XTREAM_BATCH: &str = r" -#name;username;password;url;max_connections -input_1;de566567;de2345f43g5;http://provider_1.tv:80;1 -input_2;de566567;de2345f43g5;http://provider_2.tv:8080;1 +#name;username;password;url;max_connections;exp_date +input_1;de566567;de2345f43g5;http://provider_1.tv:80;1;2028-11-23 13:12:34 +input_2;de566567;de2345f43g5;http://provider_2.tv:8080;1;2028-12-23 13:12:34 "; #[test] diff --git a/backend/src/utils/network/xtream.rs b/backend/src/utils/network/xtream.rs index 75ea86b1e..8b394d5d9 100644 --- a/backend/src/utils/network/xtream.rs +++ b/backend/src/utils/network/xtream.rs @@ -1,20 +1,21 @@ use crate::api::model::AppState; use crate::messaging::send_message; -use crate::model::{Config, ConfigInput, ConfigTarget}; +use crate::model::{is_input_expired, Config, ConfigInput, ConfigTarget, XtreamLoginInfo}; use crate::model::{InputSource, ProxyUserCredentials}; use crate::processing::parser::xtream; use crate::repository::xtream_repository; use crate::repository::xtream_repository::{rewrite_xtream_series_info_content, rewrite_xtream_vod_info_content, xtream_get_input_info}; use crate::utils::request; +use chrono::{DateTime, Utc}; use log::{error, info, warn}; use shared::error::{str_to_io_error, TuliproxError}; use shared::model::{MsgKind, PlaylistEntry, PlaylistGroup, ProxyUserStatus, XtreamCluster, XtreamPlaylistItem}; use shared::utils::{extract_extension_from_url, get_i64_from_serde_value, get_string_from_serde_value}; -// use std::cmp::Ordering; use std::io::Error; use std::str::FromStr; use std::sync::Arc; -use std::time::{SystemTime, UNIX_EPOCH}; + +const THREE_DAYS_IN_SECS: i64 = 3 * 24 * 60 * 60; #[inline] pub fn get_xtream_stream_url_base(url: &str, username: &str, password: &str) -> String { @@ -134,7 +135,7 @@ const ACTIONS: [(XtreamCluster, &str, &str); 3] = [ (XtreamCluster::Video, crate::model::XC_ACTION_GET_VOD_CATEGORIES, crate::model::XC_ACTION_GET_VOD_STREAMS), (XtreamCluster::Series, crate::model::XC_ACTION_GET_SERIES_CATEGORIES, crate::model::XC_ACTION_GET_SERIES)]; -async fn xtream_login(cfg: &Config, client: &Arc, input: &InputSource, username: &str) -> Result<(), TuliproxError> { +async fn xtream_login(cfg: &Config, client: &Arc, input: &InputSource, username: &str) -> Result, TuliproxError> { let content = if let Ok(content) = request::get_input_json_content(Arc::clone(client), None, input, None).await { content } else { @@ -148,47 +149,57 @@ async fn xtream_login(cfg: &Config, client: &Arc, input: &Input } }; - match content.get("user_info") { - None => {} - Some(value) => { - if let Some(status_value) = value.get("status") { - if let Some(status) = get_string_from_serde_value(status_value) { - if let Ok(cur_status) = ProxyUserStatus::from_str(&status) { - if !matches!(cur_status, ProxyUserStatus::Active | ProxyUserStatus::Trial) { - warn!("User status for user {username} is {cur_status:?}"); - send_message(client, MsgKind::Info, cfg.messaging.as_ref(), &format!("User status for user {username} is {cur_status:?}")).await; - } - } - } - } - if let Some(status_value) = value.get("exp_date") { - if let Some(expiration_timestamp) = get_i64_from_serde_value(status_value) { - if expiration_timestamp > 0 { - #[allow(clippy::cast_sign_loss)] - let expiration_ts = expiration_timestamp as u64; - if let Ok(now) = SystemTime::now().duration_since(UNIX_EPOCH) { - let now_secs = now.as_secs(); - if expiration_ts > now_secs { - let time_left = expiration_ts - now_secs; - if time_left < 3 * 24 * 60 * 60 { - if let Some(datetime) = chrono::DateTime::from_timestamp(expiration_timestamp, 0) { - let formatted = datetime.format("%Y-%m-%d %H:%M:%S").to_string(); - warn!("User account for user {username} expires {formatted}"); - send_message(client, MsgKind::Info, cfg.messaging.as_ref(), &format!("User account for user {username} expires {formatted}")).await; - } - } - } else { - warn!("User account for user {username} is expired"); - send_message(client, MsgKind::Info, cfg.messaging.as_ref(), &format!("User account for user {username} is expired")).await; - } - } + let mut login_info = XtreamLoginInfo { + status: None, + exp_date: None, + }; + + if let Some(user_info) = content.get("user_info") { + if let Some(status_value) = user_info.get("status") { + if let Some(status) = get_string_from_serde_value(status_value) { + if let Ok(cur_status) = ProxyUserStatus::from_str(&status) { + login_info.status = Some(cur_status); + if !matches!(cur_status, ProxyUserStatus::Active | ProxyUserStatus::Trial) { + warn!("User status for user {username} is {cur_status:?}"); + send_message(client, MsgKind::Info, cfg.messaging.as_ref(), &format!("User status for user {username} is {cur_status:?}")).await; } } } } + + if let Some(exp_value) = user_info.get("exp_date") { + if let Some(expiration_timestamp) = get_i64_from_serde_value(exp_value) { + login_info.exp_date = Some(expiration_timestamp); + notify_account_expire(login_info.exp_date, cfg, client, username).await; + } + } } - Ok(()) + if login_info.exp_date.is_none() && login_info.status.is_none() { + Ok(None) + } else { + Ok(Some(login_info)) + } +} + +pub async fn notify_account_expire(exp_date: Option, cfg: &Config, client: &Arc, username: &str) { + if let Some(expiration_timestamp) = exp_date { + let now_secs = Utc::now().timestamp(); // UTC-Time + if expiration_timestamp > now_secs { + let time_left = expiration_timestamp - now_secs; + + if time_left < THREE_DAYS_IN_SECS { + if let Some(datetime) = DateTime::::from_timestamp(expiration_timestamp, 0) { + let formatted = datetime.format("%Y-%m-%d %H:%M:%S").to_string(); + warn!("User account for user {username} expires {formatted}"); + send_message(client, MsgKind::Info, cfg.messaging.as_ref(), &format!("User account for user {username} expires {formatted}")).await; + } + } + } else { + warn!("User account for user {username} is expired"); + send_message(client, MsgKind::Info, cfg.messaging.as_ref(), &format!("User account for user {username} is expired")).await; + } + } } pub async fn get_xtream_playlist(cfg: &Arc, client: &Arc, input: &Arc, working_dir: &str) -> (Vec, Vec) { @@ -204,7 +215,9 @@ pub async fn get_xtream_playlist(cfg: &Arc, client: &Arc, client: &Arc, client: &Arc, input: &Arc) { - let cfg = Arc::clone(cfg); - let client = Arc::clone(client); - let input = Arc::clone(input); - tokio::spawn(async move { - if let Some(aliases) = input.aliases.as_ref() { - for alias in aliases { - // Random wait time 5–20 seconds to avoid provider block - let delay = u64::from(fastrand::u32(5..=20)); - tokio::time::sleep(tokio::time::Duration::from_secs(delay)).await; - - if let (Some(username), Some(password)) = - (alias.username.as_ref(), alias.password.as_ref()) - { - let mut input_source: InputSource = input.as_ref().into(); - input_source.username.clone_from(&alias.username); - input_source.password.clone_from(&alias.password); - input_source.url.clone_from(&alias.url); - let base_url = get_xtream_stream_url_base( - &input_source.url, - username, - password, - ); - let input_source_login = input_source.with_url(base_url.clone()); - - if let Err(err) = - xtream_login(&cfg, &client, &input_source_login, username).await - { - error!( - "Could not log in with xtream user {} for provider {}. {err}", - username, - alias.name - ); - } - } +async fn check_alias_user_state(cfg: &Arc, client: &Arc, input: &Arc) { + if let Some(aliases) = input.aliases.as_ref() { + for alias in aliases { + if is_input_expired(alias.exp_date) { + notify_account_expire(alias.exp_date, cfg, client, alias.username.as_ref().map_or("", |s| s.as_str())).await; } } - }); + } + + // TODO figure out how and when to call it to avoid provider bans. Possible reason for provider ban is to avoid brute force attacks. + + // + // let cfg = Arc::clone(cfg); + // let client = Arc::clone(client); + // let input = Arc::clone(input); + // + // tokio::spawn(async move { + // for alias in &aliases { + // // Random wait time 60–180 seconds to avoid provider block + // let delay = u64::from(fastrand::u32(60..=180)); + // tokio::time::sleep(tokio::time::Duration::from_secs(delay)).await; + // + // if let (Some(username), Some(password)) = + // (alias.username.as_ref(), alias.password.as_ref()) + // { + // let mut input_source: InputSource = input.as_ref().into(); + // input_source.username.clone_from(&alias.username); + // input_source.password.clone_from(&alias.password); + // input_source.url.clone_from(&alias.url); + // let base_url = get_xtream_stream_url_base( + // &input_source.url, + // username, + // password, + // ); + // let input_source_login = input_source.with_url(base_url.clone()); + // + // match xtream_login(&cfg, &client, &input_source_login, username).await { + // Ok(Some(xtream_login_info)) => { + // // TODO need to update the alias + // + // } + // Ok(None) => error!("Could log in with xtream user {} for provider {}. But could not extract account info", username, alias.name), + // Err(err) => error!("Could not log in with xtream user {} for provider {}. {err}",username,alias.name), + // } + // } + // } + // }); } pub fn create_vod_info_from_item(target: &ConfigTarget, user: &ProxyUserCredentials, pli: &XtreamPlaylistItem, last_updated: i64) -> String { @@ -315,4 +335,4 @@ pub fn create_vod_info_from_item(target: &ConfigTarget, user: &ProxyUserCredenti "stream_id": {stream_id} }} }}"#) -} \ No newline at end of file +} diff --git a/frontend/Cargo.toml b/frontend/Cargo.toml index 04d3b4292..e4a35ac45 100644 --- a/frontend/Cargo.toml +++ b/frontend/Cargo.toml @@ -1,10 +1,10 @@ [package] name = "frontend" -version = "3.2.10" +version = "3.2.11" edition = "2021" [dependencies] -shared = { version = "3.2.10", path = "../shared" } +shared = { version = "3.2.11", path = "../shared" } chrono = "0" yew = "0.21" yew-router = "0.18" diff --git a/frontend/public/assets/i18n/en.json b/frontend/public/assets/i18n/en.json index 1c32f8fd6..bf9caf1ab 100644 --- a/frontend/public/assets/i18n/en.json +++ b/frontend/public/assets/i18n/en.json @@ -8,6 +8,7 @@ "LIVE_SHORT": "L", "VOD_SHORT": "V", "SERIES_SHORT": "S", + "MOVIE": "Movie", "SAVE": "Save", "SUBMIT": "Submit", "OK": "Ok", @@ -327,7 +328,7 @@ "COPY_CREDENTIALS": "Copy Credentials" }, "TITLE": { - "USER_BOUQUET_EDITOR": "User group editor" + "USER_BOUQUET_EDITOR": "Playlist Category Selection" }, "MESSAGES": { "NO_CONTENT": "No content", @@ -341,6 +342,7 @@ "SCHEDULE_EXISTS": "Schedule already exists", "CLIPBOARD_NOT_SUPPORTED": "Clipboard not supported.\nYour browser or current context does not allow clipboard access.\nPlease use HTTPS or localhost.", "FAILED_TO_KICK_USER_STREAM": "Failed to kick user stream", + "FAILED_TO_RETRIEVE_WEBPLAYER_URL": "Failed to retrieve webplayer URL", "DOWNLOAD": { "SUCCESS": "Successfully downloaded", "FAIL": "Failed to download!", @@ -354,6 +356,9 @@ "GEOIP": { "SUCCESS": "Successfully downloaded Geo-IP db", "FAIL": "Failed to download Geo-IP db!" + }, + "USER_BOUQUET": { + "FAIL": "Failed to download user bouquets!" } }, "LOGIN": { diff --git a/frontend/public/assets/icons.json b/frontend/public/assets/icons.json index eef6df83d..4f0bd687f 100644 --- a/frontend/public/assets/icons.json +++ b/frontend/public/assets/icons.json @@ -201,7 +201,13 @@ "keys": [ "Checked" ], - "path": "M 4.2222224,2 C 3,2 2,3.000021 2,4.2222362 V 19.777764 C 2,20.999979 3,22 4.2222224,22 H 19.777778 C 20.999989,22 22,20.999979 22,19.777764 V 4.2222362 C 22,3.000021 20.999989,2 19.777778,2 Z m 0,2.2222362 H 19.777778 V 19.777764 H 4.2222224 Z M 17.288622,7.1063517 10.441833,13.953134 6.7113668,10.233533 5.2465223,11.698352 10.441833,16.89369 18.753467,8.5820472 Z" + "path": "M12 7c-2.76 0-5 2.24-5 5s2.24 5 5 5 5-2.24 5-5-2.24-5-5-5zm0-5C6.48 2 2 6.48 2 12s4.48 10 10 10 10-4.48 10-10S17.52 2 12 2zm0 18c-4.42 0-8-3.58-8-8s3.58-8 8-8 8 3.58 8 8-3.58 8-8 8z" + }, + { + "keys": [ + "Unchecked" + ], + "path": "M12 2C6.48 2 2 6.48 2 12s4.48 10 10 10 10-4.48 10-10S17.52 2 12 2zm0 18c-4.42 0-8-3.58-8-8s3.58-8 8-8 8 3.58 8 8-3.58 8-8 8z" }, { "keys": [ diff --git a/frontend/scss/_size.scss b/frontend/scss/_size.scss index aa087a6b0..bf8c9a1e5 100644 --- a/frontend/scss/_size.scss +++ b/frontend/scss/_size.scss @@ -17,6 +17,8 @@ $form-grid-cell-width-small-screen: 200px; $form-field-min-width-small-screen: 160px; $form-max-grid-cells-small-screen: 2; +$playlist-categories-grid-cell-width: 400px; + @mixin fixed-width($width) { width: $width; min-width: $width; diff --git a/frontend/scss/app/components/api_user/_api_user_view.scss b/frontend/scss/app/components/api_user/_api_user_view.scss index f62e230b6..fde91d0ac 100644 --- a/frontend/scss/app/components/api_user/_api_user_view.scss +++ b/frontend/scss/app/components/api_user/_api_user_view.scss @@ -1,3 +1,136 @@ -div { - color: var(--text-color); +@use "../../../size" as size; + +.tp__api-user-playlist { + + display: flex; + flex-flow: column; + gap: var(--gap-default); + box-sizing: border-box; + overflow: hidden; + padding: var(--padding-small); + width: 100%; + + &__loading { + transform: translateY(-10px); + } + + &__header { + flex-flow: row wrap; + gap: var(--gap-default); + + &-toolbar { + display: flex; + flex-flow: row wrap; + gap: var(--gap-default); + box-sizing: border-box; + } + } + + &__content { + display: flex; + flex-flow: column; + gap: var(--gap-default); + box-sizing: border-box; + overflow: hidden; + width: 100%; + + &-toolbar { + display: flex; + flex-flow: row; + box-sizing: border-box; + width: 100%; + justify-content: space-between; + } + + &-panels { + display: flex; + flex-flow: column; + width: 100%; + gap: var(--gap-default); + box-sizing: border-box; + overflow: hidden; + } + } +} + +.tp__api-user-target-playlist { + display: flex; + flex-flow: column; + width: 100%; + gap: var(--gap-default); + box-sizing: border-box; + overflow: hidden; + + &__body { + display: flex; + flex-flow: column; + width: 100%; + gap: var(--gap-default); + padding: var(--padding-default) 0; + box-sizing: border-box; + overflow: auto; + + .tp__collapse-panel__header { + font-size: 1.6rem; + } + + > :nth-child(1) { + border-left: 2px solid var(--output-m3u-color); + } + + > :nth-child(2) { + border-left: 2px solid var(--output-xtream-color); + } + + > :nth-child(3) { + border-left: 2px solid var(--output-hdhomerun-color); + } + + } + + &__categories { + display: grid; + grid-template-columns: repeat(auto-fit, minmax(size.$playlist-categories-grid-cell-width, 1fr)); + gap: var(--gap-large); + padding: var(--padding-default) 0; + box-sizing: border-box; + width: 100%; + + &-category.selected { + background-color: var(--text-button-active-background-color); + color: var(--text-button-active-color); + fill: var(--text-button-active-color); + } + + &-category { + display: flex; + justify-content: flex-start; + align-items: center; + background-color: var(--text-button-background-color); + color: var(--text-button-color); + fill: var(--text-button-color); + border: 1px solid var(--text-button-border-color); + border-radius: var(--border-radius); + cursor: pointer; + box-sizing: border-box; + max-height: 3rem; + min-height: 2.5rem; + gap: var(--gap-default); + padding: 0 var(--padding-default); + overflow: hidden; + text-overflow: ellipsis; + + svg, img { + pointer-events: none; + height: 1.2rem; + width: 1.2rem; + } + + &:hover { + background-color: var(--text-button-hover-background-color); + color: var(--text-button-hover-color); + } + } + } + } \ No newline at end of file diff --git a/frontend/src/app/components/api_user/api_user_view.rs b/frontend/src/app/components/api_user/api_user_view.rs index 97f512431..3a798a7a0 100644 --- a/frontend/src/app/components/api_user/api_user_view.rs +++ b/frontend/src/app/components/api_user/api_user_view.rs @@ -5,7 +5,7 @@ use crate::app::components::theme::Theme; use crate::hooks::use_service_context; use crate::provider::DialogProvider; use yew::use_state; - +use crate::app::components::api_user::playlist::ApiUserPlaylist; #[function_component] pub fn ApiUserView() -> Html { @@ -50,7 +50,7 @@ pub fn ApiUserView() -> Html {
- {" TODO "} +
diff --git a/frontend/src/app/components/api_user/mod.rs b/frontend/src/app/components/api_user/mod.rs index 13e695953..6553e0328 100644 --- a/frontend/src/app/components/api_user/mod.rs +++ b/frontend/src/app/components/api_user/mod.rs @@ -1,3 +1,5 @@ mod api_user_view; +mod playlist; +mod target_playlist; pub use api_user_view::*; \ No newline at end of file diff --git a/frontend/src/app/components/api_user/playlist.rs b/frontend/src/app/components/api_user/playlist.rs new file mode 100644 index 000000000..118550665 --- /dev/null +++ b/frontend/src/app/components/api_user/playlist.rs @@ -0,0 +1,226 @@ +use crate::app::components::api_user::target_playlist::{BouquetSelection, UserTargetPlaylist}; +use crate::app::components::{Panel, RadioButtonGroup, TextButton}; +use crate::hooks::use_service_context; +use crate::model::{BusyStatus, EventMessage}; +use shared::error::TuliproxError; +use shared::info_err; +use shared::model::{PlaylistBouquetDto, PlaylistCategoriesDto, PlaylistClusterBouquetDto}; +use std::cell::RefCell; +use std::collections::HashMap; +use std::fmt; +use std::rc::Rc; +use std::str::FromStr; +use yew::prelude::*; +use yew_i18n::use_translation; + +#[derive(Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)] +enum ApiUserPlaylistPage { + Xtream, + M3u, +} + +impl FromStr for ApiUserPlaylistPage { + type Err = TuliproxError; + + fn from_str(s: &str) -> Result { + match s.to_lowercase().as_str() { + "xtream" => Ok(ApiUserPlaylistPage::Xtream), + "m3u" => Ok(ApiUserPlaylistPage::M3u), + _ => Err(info_err!(format!("Unknown api user playlist type: {s}"))), + } + } +} + +impl fmt::Display for ApiUserPlaylistPage { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + let s = match self { + ApiUserPlaylistPage::Xtream => "xtream", + ApiUserPlaylistPage::M3u => "m3u", + }; + write!(f, "{s}") + } +} + +fn to_playlist_cluster(count: (usize, usize, usize), bouquet: Option<&Rc>>) -> Option { + if let Some(bouq) = bouquet { + let selections = bouq.borrow(); + + let selected_vec = |map: &HashMap| { + let v: Vec = map.iter() + .filter(|(_, &selected)| selected) + .map(|(c, _)| c.clone()) + .collect(); + if v.is_empty() { None } else { Some(v) } + }; + + let live = selected_vec(&selections.live).filter(|v| v.len() != count.0); + let vod = selected_vec(&selections.vod).filter(|v| v.len() != count.1); + let series = selected_vec(&selections.series).filter(|v| v.len() != count.2); + + // if all three are None, return None + if live.is_none() && vod.is_none() && series.is_none() { + None + } else { + Some(PlaylistClusterBouquetDto { live, vod, series }) + } + } else { + None + } +} + + +#[function_component] +pub fn ApiUserPlaylist() -> Html { + let translate = use_translation(); + let service_ctx = use_service_context(); + let categories = use_state(|| None as Option>); + let bouquets = use_state(|| None as Option>); + let active_tab = use_state(|| ApiUserPlaylistPage::Xtream); + let playlist_types = use_memo((), |_| { + [ApiUserPlaylistPage::Xtream, ApiUserPlaylistPage::M3u].iter().map(ToString::to_string).collect::>() + }); + + // Selection reference + let selections = use_mut_ref(|| { + HashMap::>>::new() + }); + + let handle_tab_select = { + let active_tab_clone = active_tab.clone(); + Callback::from(move |page_selection: Rc>| { + if let Some(page_selection_str) = page_selection.first() { + if let Ok(page) = ApiUserPlaylistPage::from_str(page_selection_str) { + active_tab_clone.set(page) + } + } + }) + }; + + { + // ----- Load data on mount ----- + let categories = categories.clone(); + let bouquets = bouquets.clone(); + let services = service_ctx.clone(); + let translate = translate.clone(); + + use_effect_with((), move |_| { + wasm_bindgen_futures::spawn_local(async move { + services.event.broadcast(EventMessage::Busy(BusyStatus::Show)); + let result = (services.user_api.get_playlist_bouquet().await, services.user_api.get_playlist_categories().await); + match result { + (Ok(bouquet), Ok(cats)) => { + bouquets.set(bouquet.clone()); + categories.set(cats.clone()); + } + (Err(e1), Err(e2)) => { + log::error!("Failed to load bouquet: {e1:?}, categories: {e2:?}"); + services.toastr.error(translate.t("MESSAGES.DOWNLOAD.USER_BOUQUET.FAIL")); + } + (Err(e), _) | (_, Err(e)) => { + log::error!("Failed to load user data: {e:?}"); + services.toastr.error(translate.t("MESSAGES.DOWNLOAD.USER_BOUQUET.FAIL")); + } + } + + services.event.broadcast(EventMessage::Busy(BusyStatus::Hide)); + }); + || {} + }); + } + + // ----- Save handler ----- + let on_save = { + let selections = selections.clone(); + let services = service_ctx.clone(); + let translate = translate.clone(); + let categories = categories.clone(); + + Callback::from(move |_| { + let selections = selections.clone(); + let services = services.clone(); + let translate = translate.clone(); + let categories_xtream_count = categories.as_ref().and_then(|plc| plc.xtream.as_ref().map(|x| + (x.live.as_ref().map(|v| v.len()).unwrap_or(0), + x.vod.as_ref().map(|v| v.len()).unwrap_or(0), + x.series.as_ref().map(|v| v.len()).unwrap_or(0)) + )).unwrap_or((0, 0, 0)); + let categories_m3u_count = categories.as_ref().and_then(|plc| plc.m3u.as_ref().map(|x| + (x.live.as_ref().map(|v| v.len()).unwrap_or(0), + x.vod.as_ref().map(|v| v.len()).unwrap_or(0), + x.series.as_ref().map(|v| v.len()).unwrap_or(0)) + )).unwrap_or((0, 0, 0)); + + wasm_bindgen_futures::spawn_local(async move { + services.event.broadcast(EventMessage::Busy(BusyStatus::Show)); + let result = { + let selects = selections.borrow(); + PlaylistBouquetDto { + xtream: to_playlist_cluster(categories_xtream_count, selects.get(&ApiUserPlaylistPage::Xtream)), + m3u: to_playlist_cluster(categories_m3u_count, selects.get(&ApiUserPlaylistPage::M3u)), + } + }; + + match services.user_api.save_playlist_bouquet(&result).await { + Ok(()) => services.toastr.success(translate.t("MESSAGES.SAVE.BOUQUET.SUCCESS")), + Err(_) => services.toastr.error(translate.t("MESSAGES.SAVE.BOUQUET.FAIL")), + } + + services.event.broadcast(EventMessage::Busy(BusyStatus::Hide)); + }); + }) + }; + + let handle_m3u_change = { + let selections = selections.clone(); + Callback::from(move |selection: Rc>| { + selections.borrow_mut().insert(ApiUserPlaylistPage::M3u, selection); + }) + }; + + let handle_xtream_change = { + let selections = selections.clone(); + Callback::from(move |selection: Rc>| { + selections.borrow_mut().insert(ApiUserPlaylistPage::Xtream, selection); + }) + }; + + html! { +
+
+

{translate.t("TITLE.USER_BOUQUET_EDITOR") }

+
+ + +
+
+ +
+
+ +
+ +
+ + + + + + +
+
+
+ } +} diff --git a/frontend/src/app/components/api_user/target_playlist.rs b/frontend/src/app/components/api_user/target_playlist.rs new file mode 100644 index 000000000..76e1914b9 --- /dev/null +++ b/frontend/src/app/components/api_user/target_playlist.rs @@ -0,0 +1,177 @@ +use std::cell::RefCell; +use crate::app::components::{AppIcon, Card, CollapsePanel}; +use shared::model::{PlaylistClusterBouquetDto, PlaylistClusterCategoriesDto, XtreamCluster}; +use std::collections::HashMap; +use std::rc::Rc; +use std::str::FromStr; +use wasm_bindgen::JsCast; +use yew::prelude::*; +use yew_i18n::use_translation; +use crate::html_if; + +fn normalize(s: &str) -> String { + let cleaned: String = s + .chars() + .filter(|c| c.is_alphanumeric() || c.is_whitespace()) + .collect(); + cleaned.trim().to_lowercase() +} + +fn sort_opt_vec(v: &mut Option>) { + if let Some(ref mut inner) = v { + inner.sort_by_key(|a| normalize(a)); + } +} + +macro_rules! create_selection { + ($bouquet:expr, $categories:expr, $selections:expr, $field: ident) => { + if let Some(selects) = $bouquet.$field.as_ref() { + for b in selects { + $selections.$field.insert(b.clone(), true); + } + } else { + if let Some(cats) = $categories.$field.as_ref() { + for c in cats { + $selections.$field.insert(c.clone(), true); + } + } + } + }; +} + +#[derive(Clone, PartialEq, Default)] +pub struct BouquetSelection { + pub live: HashMap, + pub vod: HashMap, + pub series: HashMap, +} + +#[derive(Properties, PartialEq)] +pub struct UserTargetPlaylistProps { + pub categories: Option, + pub bouquet: Option, + pub on_change: Callback>>, +} + +#[function_component] +pub fn UserTargetPlaylist(props: &UserTargetPlaylistProps) -> Html { + let translate = use_translation(); + let bouquet_selection = use_mut_ref(BouquetSelection::default); + let playlist_categories = use_state(PlaylistClusterCategoriesDto::default); + let force_update = use_state(|| 0); + + { + let bouquet_selection = bouquet_selection.clone(); + let playlist_categories = playlist_categories.clone(); + let in_cats = props.categories.clone(); + let in_bouquet = props.bouquet.clone(); + let force_update = force_update.clone(); + use_effect_with((in_cats, in_bouquet), move |(maybe_categories, maybe_bouquet)| { + let mut selections = BouquetSelection::default(); + if let Some(categories) = maybe_categories.as_ref() { + if let Some(bouquet) = maybe_bouquet.as_ref() { + create_selection!(bouquet, categories, selections, live); + create_selection!(bouquet, categories, selections, vod); + create_selection!(bouquet, categories, selections, series); + } else { + if let Some(cats) = categories.live.as_ref() { + for c in cats { + selections.live.insert(c.clone(), true); + } + } + if let Some(cats) = categories.vod.as_ref() { + for c in cats { + selections.vod.insert(c.clone(), true); + } + } + if let Some(cats) = categories.series.as_ref() { + for c in cats { + selections.series.insert(c.clone(), true); + } + } + } + *bouquet_selection.borrow_mut() = selections; + let mut new_categories = categories.clone(); + sort_opt_vec(&mut new_categories.live); + sort_opt_vec(&mut new_categories.vod); + sort_opt_vec(&mut new_categories.series); + playlist_categories.set(new_categories); + force_update.set(*force_update + 1); + } + }); + } + + let handle_category_click = { + let on_change = props.on_change.clone(); + let bouquet_selection = bouquet_selection.clone(); + let force_update = force_update.clone(); + Callback::from(move |e: MouseEvent| { + e.prevent_default(); + if let Some(target) = e.target() { + if let Ok(element) = target.dyn_into::() { + if let Some(cluster) = element.get_attribute("data-cluster") { + if let Ok(cluster) = XtreamCluster::from_str(cluster.as_str()) { + if let Some(category) = element.get_attribute("data-category") { + let mut selections = bouquet_selection.borrow_mut(); + match cluster { + XtreamCluster::Live => { + let selected = *selections.live.get(&category).unwrap_or(&false); + selections.live.insert(category, !selected); + } + XtreamCluster::Video => { + let selected = *selections.vod.get(&category).unwrap_or(&false); + selections.vod.insert(category, !selected); + } + XtreamCluster::Series => { + let selected = *selections.series.get(&category).unwrap_or(&false); + selections.series.insert(category, !selected); + } + } + on_change.emit(bouquet_selection.clone()); + force_update.set(*force_update + 1); + } + } + } + } + } + }) + }; + + let render_category_cluster = |cluster: XtreamCluster, cats: Option<&Vec>, selections: &HashMap| { + if let Some(c) = cats { + html_if!(!c.is_empty(), { + + "LABEL.LIVE", + XtreamCluster::Video => "LABEL.MOVIE", + XtreamCluster::Series => "LABEL.SERIES" + })}> +
+ { for c.iter().map(|cat| { + let selected = *selections.get(cat).unwrap_or(&false); + html! { +
+ { &cat } +
+ }})} +
+
+
+ }) + } else { + html! {} + } + }; + + let selections = &*bouquet_selection.borrow(); + html! { +
+
+ { render_category_cluster(XtreamCluster::Live, playlist_categories.live.as_ref(), &selections.live) } + { render_category_cluster(XtreamCluster::Video, playlist_categories.vod.as_ref(), &selections.vod) } + { render_category_cluster(XtreamCluster::Series, playlist_categories.series.as_ref(), &selections.series) } +
+
+ } +} diff --git a/frontend/src/app/components/dashboard/streams_table.rs b/frontend/src/app/components/dashboard/streams_table.rs index d5c43174f..ba3171ff2 100644 --- a/frontend/src/app/components/dashboard/streams_table.rs +++ b/frontend/src/app/components/dashboard/streams_table.rs @@ -14,7 +14,9 @@ use std::rc::Rc; use std::str::FromStr; use wasm_bindgen::JsCast; use web_sys::Element; +use yew::platform::spawn_local; use yew::prelude::*; +use yew_hooks::use_clipboard; use yew_i18n::use_translation; const LIVE: &str = "Live"; @@ -24,6 +26,11 @@ const CATCHUP: &str = "Archive"; const HLS: &str = "HLS"; const DASH: &str = "DASH"; +const KICK: &str = "kick"; +const COPY_LINK_TULIPROX_VIRTUAL_ID: &str = "copy_link_tuliprox_virtual_id"; +const COPY_LINK_TULIPROX_WEBPLAYER_URL: &str = "copy_link_tuliprox_webplayer_url"; +const COPY_LINK_PROVIDER_URL: &str = "copy_link_provider_url"; + const HEADERS: [&str; 12] = [ "EMPTY", "USERNAME", @@ -70,7 +77,8 @@ pub struct StreamsTableProps { #[function_component] pub fn StreamsTable(props: &StreamsTableProps) -> Html { let translate = use_translation(); - let services = use_service_context(); + let service_ctx = use_service_context(); + let clipboard = use_clipboard(); let config_ctx = use_context::().expect("Config context not found"); let popup_anchor_ref = use_state(|| None::); let popup_is_open = use_state(|| false); @@ -214,22 +222,62 @@ pub fn StreamsTable(props: &StreamsTableProps) -> Html { }) }; + let copy_to_clipboard: Callback = { + let clipboard = clipboard.clone(); + let services = service_ctx.clone(); + let translate = translate.clone(); + Callback::from(move |text: String| { + if *clipboard.is_supported { + clipboard.write_text(text); + } else { + services.toastr.error(translate.t("MESSAGES.CLIPBOARD_NOT_SUPPORTED")); + } + }) + }; let handle_menu_click = { let popup_is_open_state = popup_is_open.clone(); let translate = translate.clone(); - let services_ctx = services.clone(); + let services = service_ctx.clone(); let selected_dto = selected_dto.clone(); + let copy_to_clipboard = copy_to_clipboard.clone(); Callback::from(move |(name, _): (String, _)| { if let Ok(action) = StreamsTableAction::from_str(&name) { match action { StreamsTableAction::Kick => { if let Some(dto) = (*selected_dto).as_ref() { - if !services_ctx.websocket.send_message(ProtocolMessage::UserAction(UserCommand::Kick(dto.addr))) { - services_ctx.toastr.error(translate.t("MESSAGES.FAILED_TO_KICK_USER_STREAM")); + if !services.websocket.send_message(ProtocolMessage::UserAction(UserCommand::Kick(dto.addr))) { + services.toastr.error(translate.t("MESSAGES.FAILED_TO_KICK_USER_STREAM")); } } } + StreamsTableAction::CopyLinkTuliproxVirtualId => { + if let Some(dto) = &*selected_dto { + copy_to_clipboard.emit(dto.channel.virtual_id.to_string()); + } + } + StreamsTableAction::CopyLinkProviderUrl => { + if let Some(dto) = &*selected_dto { + copy_to_clipboard.emit(dto.channel.url.clone()); + } + } + StreamsTableAction::CopyLinkTuliproxWebPlayerUrl => { + if let Some(dto) = &*selected_dto { + let target_id = dto.channel.target_id; + let virtual_id = dto.channel.virtual_id; + let cluster = dto.channel.cluster; + let services = services.clone(); + let translate = translate.clone(); + let copy_to_clipboard = copy_to_clipboard.clone(); + spawn_local(async move { + if let Some(url) = services.playlist.get_playlist_webplayer_url(target_id, virtual_id, cluster).await { + copy_to_clipboard.emit(url); + } else { + services.toastr.error(translate.t("MESSAGES.FAILED_TO_RETRIEVE_WEBPLAYER_URL")); + } + }); + } + } } } popup_is_open_state.set(false); @@ -245,6 +293,9 @@ pub fn StreamsTable(props: &StreamsTableProps) -> Html { definition={definition.clone()} /> + + + } @@ -259,12 +310,18 @@ pub fn StreamsTable(props: &StreamsTableProps) -> Html { #[derive(Debug, Clone, Eq, PartialEq)] enum StreamsTableAction { Kick, + CopyLinkTuliproxVirtualId, + CopyLinkTuliproxWebPlayerUrl, + CopyLinkProviderUrl, } impl Display for StreamsTableAction { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { write!(f, "{}", match self { - Self::Kick => "kick", + Self::Kick => KICK, + Self::CopyLinkTuliproxVirtualId => COPY_LINK_TULIPROX_VIRTUAL_ID, + Self::CopyLinkTuliproxWebPlayerUrl => COPY_LINK_TULIPROX_WEBPLAYER_URL, + Self::CopyLinkProviderUrl => COPY_LINK_PROVIDER_URL, }) } } @@ -273,8 +330,14 @@ impl FromStr for StreamsTableAction { type Err = TuliproxError; fn from_str(s: &str) -> Result { - if s.eq("kick") { + if s.eq(KICK) { Ok(Self::Kick) + } else if s.eq(COPY_LINK_TULIPROX_VIRTUAL_ID) { + Ok(Self::CopyLinkTuliproxVirtualId) + } else if s.eq(COPY_LINK_TULIPROX_WEBPLAYER_URL) { + Ok(Self::CopyLinkTuliproxWebPlayerUrl) + } else if s.eq(COPY_LINK_PROVIDER_URL) { + Ok(Self::CopyLinkProviderUrl) } else { create_tuliprox_error_result!(TuliproxErrorKind::Info, "Unknown Stream Action: {}", s) } diff --git a/frontend/src/app/components/playlist/playlist_explorer.rs b/frontend/src/app/components/playlist/playlist_explorer.rs index e1c669512..a55620c14 100644 --- a/frontend/src/app/components/playlist/playlist_explorer.rs +++ b/frontend/src/app/components/playlist/playlist_explorer.rs @@ -144,13 +144,16 @@ pub fn PlaylistExplorer() -> Html { if let Some(dto) = &*selected_channel { let copy_to_clipboard = copy_to_clipboard.clone(); let services = services.clone(); - let dto = dto.clone(); + let virtual_id = dto.virtual_id; + let cluster = dto.xtream_cluster.unwrap_or_default(); let translate_clone = translate_clone.clone(); let target_id = *target_id; spawn_local(async move { - if let Some(url) = services.playlist.get_playlist_webplayer_url(target_id, &dto).await { + if let Some(url) = services.playlist.get_playlist_webplayer_url(target_id, virtual_id, cluster).await { copy_to_clipboard.emit(url); services.toastr.success(translate_clone.t("MESSAGES.PLAYLIST.WEBPLAYER_URL_COPY_TO_CLIPBOARD")); + } else { + services.toastr.error(translate_clone.t("MESSAGES.FAILED_TO_RETRIEVE_WEBPLAYER_URL")); } }); } diff --git a/frontend/src/app/components/source_editor/output_xtream_form.rs b/frontend/src/app/components/source_editor/output_xtream_form.rs index 78172360c..fa81d3fb2 100644 --- a/frontend/src/app/components/source_editor/output_xtream_form.rs +++ b/frontend/src/app/components/source_editor/output_xtream_form.rs @@ -149,7 +149,7 @@ pub fn XtreamTargetOutputView(props: &XtreamTargetOutputViewProps) -> Html { {render_output()} - {"TODO"} + {"TODO..."} diff --git a/frontend/src/app/components/tabset.rs b/frontend/src/app/components/tabset.rs index 08a861bce..f2dcacac7 100644 --- a/frontend/src/app/components/tabset.rs +++ b/frontend/src/app/components/tabset.rs @@ -12,19 +12,6 @@ pub struct TabItem { pub inactive_class: Option, } -// impl TabItem { -// pub fn new(id: String, title: String, icon: String, children: Html) -> Self { -// Self { -// id, -// title, -// icon, -// children, -// active_class: None, -// inactive_class: None, -// } -// } -// } - #[derive(Properties, Clone, PartialEq)] pub struct TabSetProps { pub tabs: Rc>, diff --git a/frontend/src/hooks/use_service_context.rs b/frontend/src/hooks/use_service_context.rs index d0463f154..d8db70a12 100644 --- a/frontend/src/hooks/use_service_context.rs +++ b/frontend/src/hooks/use_service_context.rs @@ -1,12 +1,13 @@ use std::rc::Rc; use yew::prelude::*; use crate::model::WebConfig; -use crate::services::{AuthService, ConfigService, EventService, PlaylistService, StatusService, StreamsService, ToastrService, UserService, WebSocketService}; +use crate::services::{AuthService, ConfigService, EventService, PlaylistService, StatusService, StreamsService, ToastrService, UserApiService, UserService, WebSocketService}; pub struct Services { pub auth: Rc, pub config: Rc, pub user: Rc, + pub user_api: Rc, pub status: Rc, pub streams: Rc, pub event: Rc, @@ -25,6 +26,7 @@ impl Services { let playlist = Rc::new(PlaylistService::new()); let toastr = Rc::new(ToastrService::new()); let user = Rc::new(UserService::new(Rc::clone(&event))); + let user_api = Rc::new(UserApiService::new()); let websocket = Rc::new(WebSocketService::new(Rc::clone(&status), Rc::clone(&event))); Self { auth, @@ -34,6 +36,7 @@ impl Services { event, playlist, user, + user_api, toastr, websocket } diff --git a/frontend/src/services/mod.rs b/frontend/src/services/mod.rs index 25095df33..2e41e36bc 100644 --- a/frontend/src/services/mod.rs +++ b/frontend/src/services/mod.rs @@ -8,6 +8,7 @@ mod websocket_service; mod toastr_service; mod event_service; mod user_service; +mod user_api_service; mod streams_service; pub use self::auth_service::*; @@ -20,4 +21,5 @@ pub use self::websocket_service::*; pub use self::toastr_service::*; pub use self::event_service::*; pub use self::user_service::*; +pub use self::user_api_service::*; pub use self::streams_service::*; \ No newline at end of file diff --git a/frontend/src/services/playlist_service.rs b/frontend/src/services/playlist_service.rs index a5c327b97..9a7d00e91 100644 --- a/frontend/src/services/playlist_service.rs +++ b/frontend/src/services/playlist_service.rs @@ -1,6 +1,6 @@ use crate::services::{get_base_href, request_post, ACCEPT_PREFER_BIN}; use log::error; -use shared::model::{CommonPlaylistItem, EpgTv, PlaylistCategoriesResponse, PlaylistEpgRequest, PlaylistRequest, UiPlaylistCategories, WebplayerUrlRequest}; +use shared::model::{EpgTv, PlaylistCategoriesResponse, PlaylistEpgRequest, PlaylistRequest, UiPlaylistCategories, WebplayerUrlRequest, XtreamCluster}; use std::rc::Rc; use shared::utils::{concat_path_leading_slash}; @@ -41,11 +41,11 @@ impl PlaylistService { }) } - pub async fn get_playlist_webplayer_url(&self, target_id: u16, dto: &Rc) -> Option { + pub async fn get_playlist_webplayer_url(&self, target_id: u16, virtual_id: u32, cluster: XtreamCluster) -> Option { let request = WebplayerUrlRequest { target_id, - virtual_id: dto.virtual_id, - cluster: dto.xtream_cluster.unwrap_or_default(), + virtual_id, + cluster, }; request_post::<&WebplayerUrlRequest, String>(&self.playlist_api_webplayer_url_path, &request, None, Some("text/plain".to_string())).await.unwrap_or_else(|err| { error!("{err}"); diff --git a/frontend/src/services/user_api_service.rs b/frontend/src/services/user_api_service.rs new file mode 100644 index 000000000..1e97b1afe --- /dev/null +++ b/frontend/src/services/user_api_service.rs @@ -0,0 +1,41 @@ +use std::rc::Rc; +use log::error; +use shared::model::{PlaylistBouquetDto, PlaylistCategoriesDto}; +use shared::utils::{concat_path_leading_slash}; +use crate::error::Error; +use crate::services::{get_base_href, request_get, request_post}; + +#[derive(Debug, Default)] +pub struct UserApiService { + user_playlist_categories_path: String, + user_playlist_bouquet_path: String, +} + +impl UserApiService { + pub fn new() -> Self { + let base_href = get_base_href(); + Self { + user_playlist_categories_path: concat_path_leading_slash(&base_href, "api/v1/user/playlist/categories"), + user_playlist_bouquet_path: concat_path_leading_slash(&base_href, "api/v1/user/playlist/bouquet"), + } + } + + pub async fn get_playlist_categories(&self) -> Result>, Error> { + request_get::>(&self.user_playlist_categories_path, None, None) + .await + .inspect_err(|err| error!("{err}")) + } + + pub async fn get_playlist_bouquet(&self) -> Result>, Error> { + request_get::>(&self.user_playlist_bouquet_path, None, None) + .await + .inspect_err(|err| error!("{err}")) + } + + pub async fn save_playlist_bouquet(&self, bouquet: &PlaylistBouquetDto) -> Result<(), Error> { + request_post::<&PlaylistBouquetDto, ()>(&self.user_playlist_bouquet_path, bouquet, None, None) + .await + .inspect_err(|err| error!("{err}")) + .map(|_| ()) + } +} diff --git a/shared/Cargo.toml b/shared/Cargo.toml index 51385ebaa..c12e436a7 100644 --- a/shared/Cargo.toml +++ b/shared/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "shared" -version = "3.2.10" +version = "3.2.11" edition = "2021" [dependencies] diff --git a/shared/src/model/config/api_user.rs b/shared/src/model/config/api_user.rs index 47e6a90a0..42910a272 100644 --- a/shared/src/model/config/api_user.rs +++ b/shared/src/model/config/api_user.rs @@ -1,4 +1,4 @@ -use crate::utils::default_as_true; +use crate::utils::{default_as_true, deserialize_timestamp}; use crate::error::{TuliproxError, TuliproxErrorKind}; use crate::model::{ProxyType, ProxyUserStatus}; @@ -24,7 +24,7 @@ pub struct ProxyUserCredentialsDto { pub epg_timeshift: Option, #[serde(skip_serializing_if = "Option::is_none")] pub created_at: Option, - #[serde(skip_serializing_if = "Option::is_none")] + #[serde(default, deserialize_with = "deserialize_timestamp", skip_serializing_if = "Option::is_none")] pub exp_date: Option, #[serde(default)] pub max_connections: u32, @@ -72,8 +72,8 @@ impl ProxyUserCredentialsDto { } } if let Some(exp_date) = self.exp_date { - let now = chrono::Local::now(); - if (exp_date - now.timestamp()) < 0 { + let now = chrono::Utc::now().timestamp(); + if exp_date < now { return false; } } diff --git a/shared/src/model/config/input.rs b/shared/src/model/config/input.rs index 29310c92f..e2347715c 100644 --- a/shared/src/model/config/input.rs +++ b/shared/src/model/config/input.rs @@ -1,6 +1,6 @@ use crate::error::{TuliproxError, TuliproxErrorKind}; use crate::model::{EpgConfigDto}; -use crate::utils::{default_as_true, get_credentials_from_url_str, get_trimmed_string, sanitize_sensitive_info, trim_last_slash}; +use crate::utils::{default_as_true, get_credentials_from_url_str, get_trimmed_string, sanitize_sensitive_info, trim_last_slash, deserialize_timestamp}; use crate::{check_input_credentials, check_input_connections, create_tuliprox_error_result, handle_tuliprox_error_result_list, info_err}; use enum_iterator::Sequence; use std::collections::{HashMap, HashSet}; @@ -241,6 +241,9 @@ pub struct ConfigInputAliasDto { pub priority: i16, #[serde(default)] pub max_connections: u16, + #[serde(default, deserialize_with = "deserialize_timestamp", skip_serializing_if = "Option::is_none")] + pub exp_date: Option, + } impl ConfigInputAliasDto { @@ -295,6 +298,8 @@ pub struct ConfigInputDto { pub method: InputFetchMethod, #[serde(default, skip_serializing_if = "Option::is_none")] pub staged: Option, + #[serde(default, deserialize_with = "deserialize_timestamp", skip_serializing_if = "Option::is_none")] + pub exp_date: Option, } impl Default for ConfigInputDto { @@ -316,6 +321,7 @@ impl Default for ConfigInputDto { max_connections: 0, method: InputFetchMethod::default(), staged: None, + exp_date: None, } } } diff --git a/shared/src/model/playlist.rs b/shared/src/model/playlist.rs index 032cd5ee0..937ebcf85 100644 --- a/shared/src/model/playlist.rs +++ b/shared/src/model/playlist.rs @@ -37,6 +37,19 @@ impl XtreamCluster { } } +impl FromStr for XtreamCluster { + type Err = String; + + fn from_str(s: &str) -> Result { + match s.to_lowercase().as_str() { + "live" => Ok(XtreamCluster::Live), + "video" | "vod" | "movie" => Ok(XtreamCluster::Video), + "series" => Ok(XtreamCluster::Series), + _ => Err(format!("Invalid XtreamCluster: {s}")), + } + } +} + impl Display for XtreamCluster { fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result { write!(f, "{}", self.as_str()) diff --git a/shared/src/model/playlist_categories.rs b/shared/src/model/playlist_categories.rs index f175834cd..68d0fa1ab 100644 --- a/shared/src/model/playlist_categories.rs +++ b/shared/src/model/playlist_categories.rs @@ -1,31 +1,5 @@ - -// #[derive(Debug, serde::Serialize, serde::Deserialize, Default)] -// pub struct PlaylistCategories { -// #[serde(skip_serializing_if = "Option::is_none")] -// pub live: Option>, -// #[serde(skip_serializing_if = "Option::is_none")] -// pub vod: Option>, -// #[serde(skip_serializing_if = "Option::is_none")] -// pub series: Option>, -// } - -#[derive(Debug, serde::Serialize, serde::Deserialize, Default)] -pub struct PlaylistCategoryDto { - pub id: String, - pub name: String, -} -#[derive(Debug, serde::Serialize, serde::Deserialize, Default)] -pub struct PlaylistCategoriesDto { - #[serde(skip_serializing_if = "Option::is_none")] - pub live: Option>, - #[serde(skip_serializing_if = "Option::is_none")] - pub vod: Option>, - #[serde(skip_serializing_if = "Option::is_none")] - pub series: Option>, -} - -#[derive(Debug, serde::Serialize, serde::Deserialize, Default)] -pub struct TargetBouquetDto { +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, Default, PartialEq)] +pub struct PlaylistClusterCategoriesDto { #[serde(skip_serializing_if = "Option::is_none")] pub live: Option>, #[serde(skip_serializing_if = "Option::is_none")] @@ -34,10 +8,28 @@ pub struct TargetBouquetDto { pub series: Option>, } -#[derive(Debug, serde::Serialize, serde::Deserialize, Default)] +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, Default, PartialEq)] +pub struct PlaylistCategoriesDto { + #[serde(skip_serializing_if = "Option::is_none")] + pub xtream: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub m3u: Option, +} + +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, Default, PartialEq)] +pub struct PlaylistClusterBouquetDto { + #[serde(skip_serializing_if = "Option::is_none")] + pub live: Option>, + #[serde(skip_serializing_if = "Option::is_none")] + pub vod: Option>, + #[serde(skip_serializing_if = "Option::is_none")] + pub series: Option>, +} + +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, Default, PartialEq)] pub struct PlaylistBouquetDto { #[serde(skip_serializing_if = "Option::is_none")] - pub xtream: Option, + pub xtream: Option, #[serde(skip_serializing_if = "Option::is_none")] - pub m3u: Option, + pub m3u: Option, } \ No newline at end of file diff --git a/shared/src/model/stream_info.rs b/shared/src/model/stream_info.rs index 217239222..20d432713 100644 --- a/shared/src/model/stream_info.rs +++ b/shared/src/model/stream_info.rs @@ -5,6 +5,7 @@ use crate::utils::{current_time_secs, longest}; #[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] pub struct StreamChannel { + pub target_id: u16, pub virtual_id: u32, pub provider_id: u32, pub item_type: PlaylistItemType, @@ -15,15 +16,16 @@ pub struct StreamChannel { pub shared: bool, } -pub fn create_stream_channel_with_type(pli: &XtreamPlaylistItem, item_type: PlaylistItemType) -> StreamChannel { - let mut stream_channel = pli.to_stream_channel(); +pub fn create_stream_channel_with_type(target_id: u16, pli: &XtreamPlaylistItem, item_type: PlaylistItemType) -> StreamChannel { + let mut stream_channel = pli.to_stream_channel(target_id); stream_channel.item_type = item_type; stream_channel } impl XtreamPlaylistItem { - pub fn to_stream_channel(&self) -> StreamChannel { + pub fn to_stream_channel(&self, target_id: u16) -> StreamChannel { StreamChannel { + target_id, virtual_id: self.virtual_id, provider_id: self.provider_id, item_type: self.item_type, @@ -37,8 +39,9 @@ impl XtreamPlaylistItem { } impl M3uPlaylistItem { - pub fn to_stream_channel(&self) -> StreamChannel { + pub fn to_stream_channel(&self, target_id: u16) -> StreamChannel { StreamChannel { + target_id, virtual_id: self.virtual_id, provider_id: self.get_provider_id().unwrap_or_default(), item_type: self.item_type, diff --git a/shared/src/utils/serde_utils.rs b/shared/src/utils/serde_utils.rs index 50fe265d9..080f55f56 100644 --- a/shared/src/utils/serde_utils.rs +++ b/shared/src/utils/serde_utils.rs @@ -1,8 +1,9 @@ -use std::io; -use serde::{Deserialize, Deserializer,}; -use serde::de::DeserializeOwned; -use serde_json::Value; use crate::error::to_io_error; +use chrono::{NaiveDateTime, ParseError, TimeZone, Utc}; +use serde::de::DeserializeOwned; +use serde::Deserialize; +use serde_json::Value; +use std::io; fn value_to_string_array(value: &[Value]) -> Vec { value.iter().filter_map(value_to_string).collect() @@ -19,9 +20,9 @@ fn value_to_string(v: &Value) -> Option { pub fn deserialize_as_option_string<'de, D>(deserializer: D) -> Result, D::Error> where - D: Deserializer<'de>, + D: serde::Deserializer<'de>, { - let value: Value = Deserialize::deserialize(deserializer)?; + let value: Value = serde::Deserialize::deserialize(deserializer)?; match &value { Value::String(s) => Ok(Some(s.to_owned())), @@ -32,9 +33,9 @@ where pub fn deserialize_as_string<'de, D>(deserializer: D) -> Result where - D: Deserializer<'de>, + D: serde::Deserializer<'de>, { - let value: Value = Deserialize::deserialize(deserializer)?; + let value: Value = serde::Deserialize::deserialize(deserializer)?; match &value { Value::String(s) => Ok(s.to_string()), @@ -45,7 +46,7 @@ where pub fn deserialize_as_string_array<'de, D>(deserializer: D) -> Result>, D::Error> where - D: Deserializer<'de>, + D: serde::Deserializer<'de>, { Value::deserialize(deserializer).map(|v| match v { Value::String(value) => Some(vec![value]), @@ -54,11 +55,11 @@ where }) } -pub fn deserialize_number_from_string<'de, D, T: DeserializeOwned + std::str::FromStr>( +pub fn deserialize_number_from_string<'de, D, T: DeserializeOwned + std::str::FromStr>( deserializer: D, ) -> Result, D::Error> where - D: Deserializer<'de>, + D: serde::Deserializer<'de>, { // we define a local enum type inside of the function // because it is untagged, serde will deserialize as the first variant @@ -145,4 +146,41 @@ where S: serde::Serializer, { 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>, +{ + // - try to deserialize as seconds + // - try to deserialize as date-time string of format like "2028-11-23 14:12:34" + let val = Option::::deserialize(deserializer)?; + match val { + Some(Value::Number(n)) => n + .as_i64() + .ok_or_else(|| serde::de::Error::custom("invalid number")) + .map(Some), + Some(Value::String(s)) => parse_timestamp(&s).map_err(serde::de::Error::custom), + Some(Value::Null) => Ok(None), + Some(_) => Err(serde::de::Error::custom("expected number or string")), + None => Ok(None), + } +} + +pub fn parse_timestamp(value: &str) -> Result, ParseError> { + let value = value.trim(); + if value.is_empty() { + return Ok(None); + } + + if let Ok(ts) = value.parse::() { + return Ok(Some(ts)); + } + + // "YYYY-MM-DD HH:MM:SS" + let dt = NaiveDateTime::parse_from_str(value, "%Y-%m-%d %H:%M:%S")?; + let timestamp = Utc.from_utc_datetime(&dt).timestamp(); + Ok(Some(timestamp)) } \ No newline at end of file