diff --git a/CHANGELOG.md b/CHANGELOG.md index 86d9d9616..31e328dcc 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -12,6 +12,7 @@ Output filters are applied after all transformations have been performed, theref - Added burst buffer to shared stream - Telegram message thread support. thread id can now be appended to chat-id like `chat-id:thread-id`. - Telegram supports markdown generation for structured json messages. simply set `markdown: true` in telegram config. +- Added User-Stream-Connections Table to WebUI # 3.1.7 (2025-10-10) - Added Dark/Bright theme switch diff --git a/README.md b/README.md index 4728abce6..b011b76d0 100644 --- a/README.md +++ b/README.md @@ -775,7 +775,6 @@ Has the following top level entries: - `filter` _mandatory_, - `rename` _optional_ - `mapping` _optional_ -- `favourites` _optional_ - `watch` _optional_ - `use_memory_cache`, default is false. If set to `true` playlist is cached into memory to reduce disc access. Placing playlist into memory causes more RAM usage but reduces disk access. @@ -1113,28 +1112,6 @@ messaging: - watch ``` -### 2.5.2.9 `favourites` -The Favourites feature allows users to define `custom playlist groups` that automatically collect channels based on filter rules. -Each favourite definition specifies: - -A `group name` (the name under which the favourite channels will appear in the playlist). - -A `filter expression` (which channels should be included in that group). - -During playlist processing, the system scans all existing playlist groups and channels. -For each configured favourite, it creates a new virtual group that contains all channels matching the filter condition. - -This allows a channel to appear in multiple groups — both its original one and any favourite groups where it matches. - -Example -```yaml -favourites: - - group: "[Sports] Tennis" - filter: 'Group ~ "EU: Sports" AND Caption ~ "Davis Cup Qualifiers"' - - group: "[Sports] Soccer" - filter: 'Group ~ "EU: Sports" AND Caption ~ "FIFA World Cup"' -``` - ## 2. `mapping.yml` Has the root item `mappings` which has the following top level entries: - `templates` _optional_ diff --git a/backend/src/api/api_utils.rs b/backend/src/api/api_utils.rs index 7bfc806e0..4d7da1ef3 100644 --- a/backend/src/api/api_utils.rs +++ b/backend/src/api/api_utils.rs @@ -25,10 +25,7 @@ use futures::{StreamExt, TryStreamExt}; use jsonwebtoken::{decode, Algorithm, DecodingKey, Validation}; use log::{debug, error, log_enabled, trace}; use reqwest::header::RETRY_AFTER; -use shared::model::{ - Claims, InputFetchMethod, PlaylistEntry, PlaylistItemType, TargetType, - UserConnectionPermission, XtreamCluster, -}; +use shared::model::{Claims, InputFetchMethod, PlaylistEntry, PlaylistItemType, StreamChannel, TargetType, UserConnectionPermission, XtreamCluster}; use shared::utils::{bin_serialize, default_grace_period_millis, human_readable_byte_size, trim_slash}; use shared::utils::{ extract_extension_from_url, replace_url_extension, sanitize_sensitive_info, DASH_EXT, HLS_EXT, @@ -738,7 +735,7 @@ pub async fn force_provider_stream_response( addr: &str, app_state: &AppState, user_session: &UserSession, - item_type: PlaylistItemType, + stream_channel: StreamChannel, req_headers: &HeaderMap, input: &ConfigInput, user: &ProxyUserCredentials, @@ -746,6 +743,7 @@ pub async fn force_provider_stream_response( let stream_options = get_stream_options(app_state); let share_stream = false; let connection_permission = UserConnectionPermission::Allowed; + let item_type = stream_channel.item_type; let mut stream_details = create_stream_response_details( app_state, @@ -771,7 +769,7 @@ pub async fn force_provider_stream_response( .update_session_addr(&user.username, &user_session.token, addr) .await; let stream = - ActiveClientStream::new(stream_details, app_state, user, connection_permission, addr) + ActiveClientStream::new(stream_details, app_state, user, connection_permission, addr, stream_channel) .await; let (status_code, header_map) = @@ -809,8 +807,7 @@ pub async fn stream_response( addr: &str, app_state: &AppState, session_token: &str, - virtual_id: u32, - item_type: PlaylistItemType, + stream_channel: StreamChannel, stream_url: &str, req_headers: &HeaderMap, input: &ConfigInput, @@ -830,10 +827,13 @@ pub async fn stream_response( .into_response(); } + let virtual_id = stream_channel.virtual_id; + let item_type = stream_channel.item_type; + let share_stream = is_stream_share_enabled(item_type, target); if share_stream { if let Some(value) = - shared_stream_response(app_state, stream_url, addr, user, connection_permission).await + shared_stream_response(app_state, stream_url, addr, user, connection_permission, stream_channel.clone()).await { return value.into_response(); } @@ -870,7 +870,7 @@ pub async fn stream_response( None }; let stream = - ActiveClientStream::new(stream_details, app_state, user, connection_permission, addr) + ActiveClientStream::new(stream_details, app_state, user, connection_permission, addr, stream_channel) .await; let stream_resp = if share_stream { debug_if_enabled!( @@ -984,6 +984,7 @@ async fn shared_stream_response( addr: &str, user: &ProxyUserCredentials, connect_permission: UserConnectionPermission, + stream_channel: StreamChannel ) -> Option { if let Some(stream) = SharedStreamManager::subscribe_shared_stream(app_state, stream_url, Some(addr)).await @@ -1003,7 +1004,7 @@ async fn shared_stream_response( ))); let stream_details = StreamDetails::from_stream(stream); let stream = - ActiveClientStream::new(stream_details, app_state, user, connect_permission, addr) + ActiveClientStream::new(stream_details, app_state, user, connect_permission, addr, stream_channel) .await .boxed(); let mut response = axum::response::Response::builder().status(status_code); diff --git a/backend/src/api/endpoints/api_playlist_utils.rs b/backend/src/api/endpoints/api_playlist_utils.rs index d8af876c8..bb84329da 100644 --- a/backend/src/api/endpoints/api_playlist_utils.rs +++ b/backend/src/api/endpoints/api_playlist_utils.rs @@ -128,7 +128,7 @@ async fn grouped_channels( pub(in crate::api::endpoints) async fn get_playlist_for_target(cfg_target: Option<&ConfigTarget>, cfg: &AppConfig, accept: Option<&String>) -> impl axum::response::IntoResponse + Send { if let Some(target) = cfg_target { - if target.has_output(&TargetType::Xtream) { + if target.has_output(TargetType::Xtream) { let live_channels = grouped_channels(cfg, target, XtreamCluster::Live).await; let vod_channels = grouped_channels(cfg, target, XtreamCluster::Video).await; let series_channels = grouped_channels(cfg, target, XtreamCluster::Series).await; @@ -140,7 +140,7 @@ pub(in crate::api::endpoints) async fn get_playlist_for_target(cfg_target: Optio }; return json_or_bin_response(accept, &response).into_response(); - } else if target.has_output(&TargetType::M3u) { + } else if target.has_output(TargetType::M3u) { let all_channels = m3u_repository::iter_raw_m3u_playlist(cfg, target).await; let (live_channels, vod_channels, series_channels) = group_playlist_items_by_cluster(all_channels); let response = PlaylistCategoriesResponse { diff --git a/backend/src/api/endpoints/hdhomerun_api.rs b/backend/src/api/endpoints/hdhomerun_api.rs index c7787d99a..0136ad79a 100644 --- a/backend/src/api/endpoints/hdhomerun_api.rs +++ b/backend/src/api/endpoints/hdhomerun_api.rs @@ -224,7 +224,7 @@ async fn lineup(app_state: &Arc, cfg: &Arc, creden let use_all = use_output.is_none(); let use_m3u = use_output.as_ref() == Some(&TargetType::M3u); let use_xtream = use_output.as_ref() == Some(&TargetType::Xtream); - if (use_all || use_m3u) && target.has_output(&TargetType::M3u) { + if (use_all || use_m3u) && target.has_output(TargetType::M3u) { let iterator = M3uPlaylistIterator::new(cfg, target, credentials).await.ok(); let stream = m3u_item_to_lineup_stream(iterator); let body_stream = stream::once(async { Ok(Bytes::from("[")) }) @@ -234,7 +234,7 @@ async fn lineup(app_state: &Arc, cfg: &Arc, creden .status(axum::http::StatusCode::OK) .header(axum::http::header::CONTENT_TYPE, mime::APPLICATION_JSON.to_string()) .body(axum::body::Body::from_stream(body_stream))); - } else if (use_all || use_xtream) && target.has_output(&TargetType::Xtream) { + } else if (use_all || use_xtream) && target.has_output(TargetType::Xtream) { let server_info = app_state.app_state.app_config.get_user_server_info(credentials); let base_url = server_info.get_base_url(); diff --git a/backend/src/api/endpoints/hls_api.rs b/backend/src/api/endpoints/hls_api.rs index 85f57caa5..6af9cda86 100644 --- a/backend/src/api/endpoints/hls_api.rs +++ b/backend/src/api/endpoints/hls_api.rs @@ -7,7 +7,7 @@ use crate::api::model::UserSession; use crate::api::model::{create_custom_video_stream_response, CustomVideoStreamType}; use crate::api::model::AppState; use crate::auth::Fingerprint; -use crate::model::ProxyUserCredentials; +use crate::model::{ConfigTarget, ProxyUserCredentials}; use crate::model::{ConfigInput, InputSource}; use crate::processing::parser::hls::{ get_hls_session_token_and_url_from_token, rewrite_hls, RewriteHlsProps, @@ -17,9 +17,11 @@ use axum::http::HeaderMap; use axum::response::IntoResponse; use log::{debug, error}; use serde::Deserialize; -use shared::model::{PlaylistItemType, UserConnectionPermission, XtreamCluster}; +use shared::model::{PlaylistItemType, StreamChannel, TargetType, UserConnectionPermission, XtreamCluster}; use shared::utils::{is_hls_url, replace_url_extension, sanitize_sensitive_info, CUSTOM_VIDEO_PREFIX, HLS_EXT}; use std::sync::Arc; +use crate::repository::m3u_repository::m3u_get_item_for_stream_id; +use crate::repository::xtream_repository; #[derive(Debug, Deserialize)] struct HlsApiPathParams { @@ -151,6 +153,49 @@ 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, _mapping)) = xtream_repository::xtream_get_item_for_stream_id( + virtual_id, + app_state, + target, + None + ).await { + Some(pli.to_stream_channel()) + } else { + None + } + } else { + match m3u_get_item_for_stream_id(virtual_id, app_state, target).await { + Ok(pli) => Some(pli.to_stream_channel()), + Err(_) => { None } + } + } +} + +async fn resolve_stream_channel( + app_state: &Arc, + target: &Arc, + virtual_id: u32, + hls_url: &str, +) -> StreamChannel { + let mut channel = match get_stream_channel(app_state, target, virtual_id).await { + Some(channel) => channel, + None => StreamChannel { + virtual_id, + provider_id: 0, + item_type: PlaylistItemType::LiveHls, + cluster: XtreamCluster::Live, + group: "Unknown".to_string(), + title: "Unknown".to_string(), + url: hls_url.to_string(), + }, + }; + + channel.item_type = PlaylistItemType::LiveHls; + channel +} + #[allow(clippy::too_many_lines)] async fn hls_api_stream( Fingerprint(fingerprint, addr): Fingerprint, @@ -207,7 +252,7 @@ async fn hls_api_stream( &app_state.app_config, CustomVideoStreamType::ProviderConnectionsExhausted, ) - .into_response(); + .into_response(); } let hls_url = match get_hls_session_token_and_url_from_token( @@ -218,21 +263,22 @@ async fn hls_api_stream( _ => return axum::http::StatusCode::BAD_REQUEST.into_response(), }; - session.stream_url = hls_url; + session.stream_url.clone_from(&hls_url); if session.virtual_id == virtual_id { if is_seek_request(XtreamCluster::Live, &req_headers).await { // partial request means we are in reverse proxy mode, seek happened + let stream_channel = resolve_stream_channel(&app_state, &target, virtual_id, &hls_url).await; return force_provider_stream_response( &addr, &app_state, session, - PlaylistItemType::LiveHls, + stream_channel, &req_headers, &input, &user, ) - .await - .into_response(); + .await + .into_response(); } } else { return axum::http::StatusCode::BAD_REQUEST.into_response(); @@ -264,11 +310,13 @@ async fn hls_api_stream( .into_response(); } + let stream_channel = resolve_stream_channel(&app_state, &target, virtual_id, &hls_url).await; + force_provider_stream_response( &addr, &app_state, session, - PlaylistItemType::LiveHls, + stream_channel, &req_headers, &input, &user, @@ -280,6 +328,7 @@ async fn hls_api_stream( } } + pub fn hls_api_register() -> axum::Router> { axum::Router::new().route( "/hls/{username}/{password}/{input_id}/{stream_id}/{token}", diff --git a/backend/src/api/endpoints/m3u_api.rs b/backend/src/api/endpoints/m3u_api.rs index 27621c81b..f59836d28 100644 --- a/backend/src/api/endpoints/m3u_api.rs +++ b/backend/src/api/endpoints/m3u_api.rs @@ -102,7 +102,7 @@ async fn m3u_api_stream( } let target_name = &target.name; - if !target.has_output(&TargetType::M3u) { + if !target.has_output(TargetType::M3u) { debug!("Target has no m3u playlist {target_name}"); return axum::http::StatusCode::BAD_REQUEST.into_response(); } @@ -154,7 +154,7 @@ async fn m3u_api_stream( addr, app_state, session, - pli.item_type, + pli.to_stream_channel(), req_headers, &input, &user, @@ -188,7 +188,7 @@ async fn m3u_api_stream( user: &user, stream_ext: stream_ext.as_deref(), req_context: context, - action_path: "", // TODO is there timeshoft or something like that ? + action_path: "", // TODO is there timeshift or something like that ? }; if let Some(response) = redirect_response(app_state, &redirect_params).await { @@ -225,8 +225,7 @@ async fn m3u_api_stream( addr, app_state, &session_key, - pli.virtual_id, - pli.item_type, + pli.to_stream_channel(), session_url, req_headers, &input, @@ -262,7 +261,7 @@ async fn m3u_api_resource( } let target_name = &target.name; - if !target.has_output(&TargetType::M3u) { + if !target.has_output(TargetType::M3u) { debug!("Target has no m3u playlist {target_name}"); return axum::http::StatusCode::BAD_REQUEST.into_response(); } diff --git a/backend/src/api/endpoints/user_api.rs b/backend/src/api/endpoints/user_api.rs index f4b11195b..03a8646bc 100644 --- a/backend/src/api/endpoints/user_api.rs +++ b/backend/src/api/endpoints/user_api.rs @@ -51,7 +51,7 @@ async fn playlist_categories( return axum::http::StatusCode::FORBIDDEN.into_response(); } let target_name = &target.name; - let xtream_stream = if target.has_output(&TargetType::Xtream) { + let xtream_stream = if target.has_output(TargetType::Xtream) { let config = &app_state.app_config.config.load(); let live_categories = get_categories_from_xtream(xtream_get_playlist_categories(config, target_name, XtreamCluster::Live).await); let vod_categories = get_categories_from_xtream(xtream_get_playlist_categories(config, target_name, XtreamCluster::Video).await); @@ -69,7 +69,7 @@ async fn playlist_categories( stream::iter(vec![Ok::(Bytes::from(r#"{"live":[],"vod":[],"series":[]}"#))]) }; - let m3u_stream = if target.has_output(&TargetType::M3u) { + let m3u_stream = if target.has_output(TargetType::M3u) { let live_categories = get_categories_from_m3u_playlist(&target, &app_state.app_config).await; stream::iter(vec![ Ok::(Bytes::from(r#"{"live": "#)), diff --git a/backend/src/api/endpoints/v1_api.rs b/backend/src/api/endpoints/v1_api.rs index 4339cf709..8eddc7125 100644 --- a/backend/src/api/endpoints/v1_api.rs +++ b/backend/src/api/endpoints/v1_api.rs @@ -1,4 +1,4 @@ -use crate::api::api_utils::try_unwrap_body; +use crate::api::api_utils::{json_or_bin_response, try_unwrap_body}; use crate::api::endpoints::download_api; use crate::api::endpoints::user_api::user_api_register; use crate::api::endpoints::v1_api_playlist::v1_api_playlist_register; @@ -12,6 +12,7 @@ use shared::model::{IpCheckDto,StatusCheck}; use shared::utils::{concat_path_leading_slash}; use std::collections::BTreeMap; use std::sync::Arc; +use crate::api::endpoints::extract_accept_header::ExtractAcceptHeader; use crate::api::endpoints::v1_api_config::v1_api_config_register; async fn create_ipinfo_check(app_state: &Arc) -> Option<(Option, Option)> { @@ -31,9 +32,9 @@ pub async fn create_status_check(app_state: &Arc) -> StatusCheck { Some(lock.lock().await.get_size_text()) } }; - let (active_users, active_user_connections) = { + let (active_users, active_user_connections, active_user_streams) = { let active_user = &app_state.active_users; - (active_user.active_users().await, active_user.active_connections().await) + (active_user.active_users().await, active_user.active_connections().await, active_user.active_streams().await) }; let active_provider_connections = app_state.active_provider.active_connections().await.map(|c| c.into_iter().collect::>()); @@ -47,6 +48,7 @@ pub async fn create_status_check(app_state: &Arc) -> StatusCheck { active_users, active_user_connections, active_provider_connections, + active_user_streams, cache, } } @@ -59,6 +61,12 @@ async fn status(axum::extract::State(app_state): axum::extract::State>) -> axum::response::Response { + let streams = app_state.active_users.active_streams().await; + json_or_bin_response(accept.as_ref(), &streams).into_response() +} + async fn ipinfo(axum::extract::State(app_state): axum::extract::State>) -> axum::response::Response { if let Some((ipv4, ipv6)) = create_ipinfo_check(&app_state).await { let ipcheck = IpCheckDto { @@ -78,6 +86,7 @@ pub fn v1_api_register(web_auth_enabled: bool, app_state: Arc, web_ui_ let mut router = axum::Router::new(); router = router .route("/status", axum::routing::get(status)) + .route("/streams", axum::routing::get(streams)) .route("/file/download", axum::routing::post(download_api::queue_download_file)) .route("/file/download/info", axum::routing::get(download_api::download_file_info)) .route("/ipinfo", axum::routing::get(ipinfo)); diff --git a/backend/src/api/endpoints/websocket_api.rs b/backend/src/api/endpoints/websocket_api.rs index 0c64220ca..ced702c15 100644 --- a/backend/src/api/endpoints/websocket_api.rs +++ b/backend/src/api/endpoints/websocket_api.rs @@ -201,8 +201,8 @@ async fn handle_event_message(socket: &mut WebSocket, event: EventMessage, handl .await .map_err(|e| format!("Server Error event: {e} "))?; } - EventMessage::ActiveUser(users, connections) => { - let msg = ProtocolMessage::ActiveUserResponse(users, connections) + EventMessage::ActiveUser(event) => { + let msg = ProtocolMessage::ActiveUserResponse(event) .to_bytes() .map_err(|e| e.to_string())?; socket diff --git a/backend/src/api/endpoints/xtream_api.rs b/backend/src/api/endpoints/xtream_api.rs index 0ed14ef16..e67ef1d9f 100644 --- a/backend/src/api/endpoints/xtream_api.rs +++ b/backend/src/api/endpoints/xtream_api.rs @@ -240,7 +240,7 @@ async fn xtream_player_api_stream( } let target_name = &target.name; - if !target.has_output(&TargetType::Xtream) { + if !target.has_output(TargetType::Xtream) { debug!("Target has no xtream codes playlist {target_name}"); return axum::http::StatusCode::BAD_REQUEST.into_response(); } @@ -306,7 +306,7 @@ async fn xtream_player_api_stream( addr, app_state, session, - item_type, + pli.to_stream_channel(), req_headers, &input, &user, @@ -392,8 +392,7 @@ async fn xtream_player_api_stream( addr, app_state, session_key.as_str(), - pli.virtual_id, - item_type, + pli.to_stream_channel(), &stream_url, req_headers, &input, @@ -417,7 +416,7 @@ async fn xtream_player_api_stream_with_token( ) -> impl IntoResponse + Send { if let Some(target) = app_state.app_config.get_target_by_id(target_id) { let target_name = &target.name; - if !target.has_output(&TargetType::Xtream) { + if !target.has_output(TargetType::Xtream) { debug!("Target has no xtream output {target_name}"); return axum::http::StatusCode::BAD_REQUEST.into_response(); } @@ -524,8 +523,7 @@ async fn xtream_player_api_stream_with_token( addr, app_state, session_key.as_str(), - pli.virtual_id, - pli.item_type, + pli.to_stream_channel(), &stream_url, req_headers, &input, @@ -720,7 +718,7 @@ async fn xtream_player_api_resource( return axum::http::StatusCode::FORBIDDEN.into_response(); } let target_name = &target.name; - if !target.has_output(&TargetType::Xtream) { + if !target.has_output(TargetType::Xtream) { debug!("Target has no xtream output {target_name}"); return axum::http::StatusCode::BAD_REQUEST.into_response(); } @@ -1046,7 +1044,7 @@ async fn xtream_get_short_epg( limit: &str, ) -> impl IntoResponse + Send { let target_name = &target.name; - if target.has_output(&TargetType::Xtream) { + if target.has_output(TargetType::Xtream) { let virtual_id: u32 = match FromStr::from_str(stream_id.trim()) { Ok(id) => id, Err(_) => return axum::http::StatusCode::BAD_REQUEST.into_response(), @@ -1337,7 +1335,7 @@ async fn xtream_player_api( ) -> impl IntoResponse + Send { let user_target = get_user_target(&api_req, app_state); if let Some((user, target)) = user_target { - if !target.has_output(&TargetType::Xtream) { + if !target.has_output(TargetType::Xtream) { return axum::response::Json(get_user_info(&user, app_state).await).into_response(); } diff --git a/backend/src/api/model/active_user_manager.rs b/backend/src/api/model/active_user_manager.rs index 2ed812b10..45d26d478 100644 --- a/backend/src/api/model/active_user_manager.rs +++ b/backend/src/api/model/active_user_manager.rs @@ -4,7 +4,7 @@ use crate::model::Config; use crate::model::ProxyUserCredentials; use jsonwebtoken::get_current_timestamp; use log::{debug, error, info}; -use shared::model::UserConnectionPermission; +use shared::model::{ActiveUserConnectionChange, StreamChannel, StreamInfo, UserConnectionPermission}; use shared::utils::{current_time_secs, default_grace_period_millis, default_grace_period_timeout_secs, sanitize_sensitive_info}; use std::collections::HashMap; use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; @@ -16,8 +16,8 @@ const USER_GC_TTL: u64 = 900; // 15 Min const USER_CON_TTL: u64 = 10_800; // 3 hours const USER_SESSION_LIMIT: usize = 50; -type ActiveUserConnectionChangeSender = tokio::sync::mpsc::Sender<(usize, usize)>; -pub type ActiveUserConnectionChangeReceiver = tokio::sync::mpsc::Receiver<(usize, usize)>; +type ActiveUserConnectionChangeSender = tokio::sync::mpsc::Sender; +pub type ActiveUserConnectionChangeReceiver = tokio::sync::mpsc::Receiver; macro_rules! active_user_manager_shared_impl { () => { @@ -35,11 +35,10 @@ macro_rules! active_user_manager_shared_impl { async fn log_active_user(&self) { let user = Arc::clone(&self.user); - let connection_change_tx = self.connection_change_tx.clone(); let is_log_user_enabled = self.is_log_user_enabled(); let user_connection_count = Self::get_active_connections(&user).await; let user_count = user.read().await.iter().filter(|(_, c)| c.connections > 0).count(); - let _= connection_change_tx.try_send((user_count, user_connection_count)); + let _= self.connection_change_tx.try_send(ActiveUserConnectionChange::Connections(user_count, user_connection_count)); if is_log_user_enabled { info!("Active Users: {user_count}, Active User Connections: {user_connection_count}"); } @@ -61,11 +60,13 @@ macro_rules! active_user_manager_shared_impl { connection_data.granted_grace = false; connection_data.grace_ts = 0; } + connection_data.streams.retain(|c| c.addr != addr); } } self.drop_connection(&addr); self.shared_stream_manager.release_connection(addr, true).await; self.provider_manager.release_connection(addr).await; + let _= self.connection_change_tx.try_send(ActiveUserConnectionChange::Disconnected(addr.to_string())); self.log_active_user().await; } }; @@ -132,6 +133,7 @@ struct UserConnectionData { granted_grace: bool, grace_ts: u64, sessions: Vec, + streams: Vec, ts: u64, } @@ -143,6 +145,7 @@ impl UserConnectionData { granted_grace: false, grace_ts: 0, sessions: Vec::new(), + streams: Vec::new(), ts: current_time_secs(), } } @@ -277,16 +280,24 @@ impl ActiveUserManager { Self::get_active_connections(&self.user).await } - pub async fn add_connection(&self, username: &str, max_connections: u32, addr: &str) -> UserConnectionGuard { + pub async fn add_connection(&self, username: &str, max_connections: u32, addr: &str, provider: &str, stream_channel: StreamChannel) -> UserConnectionGuard { + let stream_info = StreamInfo::new( + username, + addr, + provider, + stream_channel, + ); { let mut user_map = self.user.write().await; - user_map - .entry(username.to_string()) - .and_modify(|connection_data| { - connection_data.connections += 1; - connection_data.max_connections = max_connections; - }) - .or_insert_with(|| UserConnectionData::new(1, max_connections)); + if let Some(connection_data) = user_map.get_mut(username) { + connection_data.connections += 1; + connection_data.max_connections = max_connections; + connection_data.streams.push(stream_info.clone()); + } else { + let mut connection_data = UserConnectionData::new(1, max_connections); + connection_data.streams.push(stream_info.clone()); + user_map.insert(username.to_string(), connection_data); + } } { @@ -294,6 +305,7 @@ impl ActiveUserManager { user_by_addr.insert(addr.to_owned(), username.to_owned()); } + let _= self.connection_change_tx.try_send(ActiveUserConnectionChange::Connected(stream_info)); self.log_active_user().await; UserConnectionGuard { @@ -351,7 +363,7 @@ impl ActiveUserManager { } } - // If not session exists, create one + // If no session exists, create one debug!("Creating session for user {} with token {session_token} {}", user.username, sanitize_sensitive_info(stream_url)); let session = Self::new_user_session(session_token, virtual_id, provider, stream_url, addr, connection_permission); let token = session.token.clone(); @@ -408,6 +420,17 @@ impl ActiveUserManager { None } + pub async fn active_streams(&self) -> Vec { + let user_map = self.user.read().await; + let mut streams = Vec::new(); + for (_username, connection_data) in user_map.iter() { + for stream in &connection_data.streams { + streams.push(stream.clone()); + } + } + streams + } + fn gc(&self) { if let Some(gc_ts) = &self.gc_ts { let ts = gc_ts.load(Ordering::Acquire); diff --git a/backend/src/api/model/event_manager.rs b/backend/src/api/model/event_manager.rs index 7efaa7565..46a738d07 100644 --- a/backend/src/api/model/event_manager.rs +++ b/backend/src/api/model/event_manager.rs @@ -1,13 +1,13 @@ use log::{info, trace}; use tokio::task; -use shared::model::{ConfigType, PlaylistUpdateState}; +use shared::model::{ActiveUserConnectionChange, ConfigType, PlaylistUpdateState}; use crate::api::model::{ActiveUserConnectionChangeReceiver}; use crate::api::model::{ProviderConnectionChangeReceiver}; #[derive(Clone, PartialEq)] pub enum EventMessage { ServerError(String), - ActiveUser(usize, usize), // user_count, connection count + ActiveUser(ActiveUserConnectionChange), // user_count, connection count ActiveProvider(String, usize), // provider name, connections ConfigChange(ConfigType), PlaylistUpdate(PlaylistUpdateState), @@ -30,22 +30,22 @@ impl EventManager { task::spawn(async move { loop { tokio::select! { - Some((user_count, connection_count)) = active_user_change_rx.recv() => { - if let Err(e) = channel_tx_clone.send(EventMessage::ActiveUser(user_count, connection_count)) { - trace!("Failed to send active user change event: {e}"); - } - } + Some(event) = active_user_change_rx.recv() => { + if let Err(e) = channel_tx_clone.send(EventMessage::ActiveUser(event)) { + trace!("Failed to send active user change event: {e}"); + } + } - Some((provider, connection_count)) = provider_change_rx.recv() => { + Some((provider, connection_count)) = provider_change_rx.recv() => { if let Err(e) = channel_tx_clone.send(EventMessage::ActiveProvider(provider, connection_count)) { trace!("Failed to send active provider change event: {e}"); } - } - else => { + } + else => { // Both channels are closed, exit gracefully info!("All input channels closed, terminating event manager task"); break; - } + } } } }); diff --git a/backend/src/api/model/streams/active_client_stream.rs b/backend/src/api/model/streams/active_client_stream.rs index a0d94810e..cb2613fd6 100644 --- a/backend/src/api/model/streams/active_client_stream.rs +++ b/backend/src/api/model/streams/active_client_stream.rs @@ -10,7 +10,7 @@ use bytes::Bytes; use futures::Stream; use futures::StreamExt; use log::{error, info}; -use shared::model::UserConnectionPermission; +use shared::model::{StreamChannel, UserConnectionPermission}; use std::pin::Pin; use std::sync::atomic::AtomicU8; use std::sync::{Arc}; @@ -38,13 +38,20 @@ impl ActiveClientStream { app_state: &AppState, user: &ProxyUserCredentials, connection_permission: UserConnectionPermission, - addr: &str) -> Self { + addr: &str, + stream_channel: StreamChannel) -> Self { if connection_permission == UserConnectionPermission::Exhausted { error!("Something is wrong this should not happen"); } let grant_user_grace_period = connection_permission == UserConnectionPermission::GracePeriod; let username = user.username.as_str(); - let user_connection_guard = Some(app_state.active_users.add_connection(username, user.max_connections, addr).await); + let provider_name = stream_details + .provider_connection_guard + .as_ref() + .and_then(|guard| guard.get_provider_name()) + .as_deref() + .map_or_else(String::new, ToString::to_string); + let user_connection_guard = Some(app_state.active_users.add_connection(username, user.max_connections, addr, &provider_name, stream_channel).await); let cfg = &app_state.app_config; let waker = Some(Arc::new(AtomicWaker::new())); let waker_clone = waker.clone(); diff --git a/backend/src/auth/fingerprint.rs b/backend/src/auth/fingerprint.rs index a401b9e53..791f09459 100644 --- a/backend/src/auth/fingerprint.rs +++ b/backend/src/auth/fingerprint.rs @@ -59,12 +59,18 @@ impl Fingerprint { } } + let client_ip = format!("{}:{}", real_ip.as_ref() + .map(ToString::to_string) + .or(forwarded_for.as_ref().map(ToString::to_string)) + .unwrap_or_else(|| addr.ip().to_string()), + addr.port()); + let ua = user_agent.unwrap_or_else(String::new); let key = match real_ip.or(forwarded_for) { Some(xff) => format!("{}{xff}{ua}", addr.ip()), None => format!("{}{ua}", addr.ip()), }; - Ok(Fingerprint(key, addr.to_string())) + Ok(Fingerprint(key, client_ip)) } } \ No newline at end of file diff --git a/backend/src/model/config/target.rs b/backend/src/model/config/target.rs index 3ce9f2ede..6401281ed 100644 --- a/backend/src/model/config/target.rs +++ b/backend/src/model/config/target.rs @@ -265,13 +265,13 @@ impl ConfigTarget { } } - pub fn has_output(&self, tt: &TargetType) -> bool { + pub fn has_output(&self, tt: TargetType) -> bool { for target_output in &self.output { match target_output { - TargetOutput::Xtream(_) => { if tt == &TargetType::Xtream { return true; } } - TargetOutput::M3u(_) => { if tt == &TargetType::M3u { return true; } } - TargetOutput::Strm(_) => { if tt == &TargetType::Strm { return true; } } - TargetOutput::HdHomeRun(_) => { if tt == &TargetType::HdHomeRun { return true; } } + TargetOutput::Xtream(_) => { if tt == TargetType::Xtream { return true; } } + TargetOutput::M3u(_) => { if tt == TargetType::M3u { return true; } } + TargetOutput::Strm(_) => { if tt == TargetType::Strm { return true; } } + TargetOutput::HdHomeRun(_) => { if tt == TargetType::HdHomeRun { return true; } } } } false diff --git a/backend/src/processing/parser/xmltv.rs b/backend/src/processing/parser/xmltv.rs index dc0a26be7..b94dcd537 100644 --- a/backend/src/processing/parser/xmltv.rs +++ b/backend/src/processing/parser/xmltv.rs @@ -393,7 +393,7 @@ fn handle_text_tag(stack: &mut [XmlTag], e: &BytesText) { let t_fixed: Cow = if t.ends_with('\\') { let mut owned = t.to_string(); owned.pop(); - owned.push_str("'"); + owned.push_str("' "); Cow::Owned(owned) } else { Cow::Borrowed(t) diff --git a/backend/src/processing/processor/playlist.rs b/backend/src/processing/processor/playlist.rs index bb45dbac5..8e36a87f7 100644 --- a/backend/src/processing/processor/playlist.rs +++ b/backend/src/processing/processor/playlist.rs @@ -61,41 +61,41 @@ pub fn apply_filter_to_playlist(playlist: &mut [PlaylistGroup], filter: &Filter) } pub fn apply_favourites_to_playlist( - playlist: &mut Vec, - favourites_cfg: Option<&Vec>, + _playlist: &mut Vec, + _favourites_cfg: Option<&Vec>, ) { - if let Some(favourites) = favourites_cfg { - let mut fav_groups: HashMap> = HashMap::new(); - - for pg in playlist.iter_mut() { - for pli in &pg.channels { - for fav in favourites { - if is_valid(pli, &fav.filter) { - let mut channel = pli.clone(); - channel.header.copy = true; - channel.header.group.clone_from(&fav.group); - channel.header.gen_uuid(); - fav_groups - .entry(fav.group.clone()) - .or_default() - .push(channel); - } - } - } - } - - for (group_name, channels) in fav_groups { - if !channels.is_empty() { - let xtream_cluster = channels[0].header.xtream_cluster; - playlist.push(PlaylistGroup { - id: 0, - title: group_name, - channels, - xtream_cluster, - }); - } - } - } + // if let Some(favourites) = favourites_cfg { + // let mut fav_groups: HashMap> = HashMap::new(); + // + // for pg in playlist.iter_mut() { + // for pli in &pg.channels { + // for fav in favourites { + // if is_valid(pli, &fav.filter) { + // let mut channel = pli.clone(); + // channel.header.copy = true; + // channel.header.group.clone_from(&fav.group); + // channel.header.gen_uuid(); + // fav_groups + // .entry(fav.group.clone()) + // .or_default() + // .push(channel); + // } + // } + // } + // } + // + // for (group_name, channels) in fav_groups { + // if !channels.is_empty() { + // let xtream_cluster = channels[0].header.xtream_cluster; + // playlist.push(PlaylistGroup { + // id: 0, + // title: group_name, + // channels, + // xtream_cluster, + // }); + // } + // } + // } } fn filter_playlist(playlist: &mut [PlaylistGroup], target: &ConfigTarget) -> Option> { diff --git a/backend/src/processing/processor/sort.rs b/backend/src/processing/processor/sort.rs index 38e34ab0b..3f2d2d26f 100644 --- a/backend/src/processing/processor/sort.rs +++ b/backend/src/processing/processor/sort.rs @@ -65,7 +65,11 @@ fn playlist_comparator( } } - Ordering::Equal + let o = value_a.cmp(value_b); + match order { + SortOrder::Asc => o, + SortOrder::Desc => o.reverse(), + } } (Some(_), None) => match order { SortOrder::Asc => Ordering::Less, diff --git a/frontend/public/assets/fonts.css b/frontend/public/assets/fonts.css index 2f6f658f8..4ce6e53dd 100644 --- a/frontend/public/assets/fonts.css +++ b/frontend/public/assets/fonts.css @@ -1,142 +1,142 @@ @font-face { - font-family: 'Montserrat Regular'; + font-family: 'Montserrat'; font-style: normal; font-weight: normal; - src: local('Montserrat Regular'), url('fonts/Montserrat-Regular.woff') format('woff'); + src: local('Montserrat'), url('fonts/montserrat-regular.woff') format('woff'); } -@font-face { - font-family: 'Montserrat Italic'; - font-style: normal; - font-weight: normal; - src: local('Montserrat Italic'), url('fonts/Montserrat-Italic.woff') format('woff'); -} +/*@font-face {*/ +/* font-family: 'Montserrat Italic';*/ +/* font-style: normal;*/ +/* font-weight: normal;*/ +/* src: local('Montserrat Italic'), url('fonts/montserrat-italic.woff') format('woff');*/ +/*}*/ -@font-face { - font-family: 'Montserrat Thin'; - font-style: normal; - font-weight: normal; - src: local('Montserrat Thin'), url('fonts/Montserrat-Thin.woff') format('woff'); -} +/*@font-face {*/ +/* font-family: 'Montserrat Thin';*/ +/* font-style: normal;*/ +/* font-weight: normal;*/ +/* src: local('Montserrat Thin'), url('fonts/montserrat-thin.woff') format('woff');*/ +/*}*/ -@font-face { - font-family: 'Montserrat Thin Italic'; - font-style: normal; - font-weight: normal; - src: local('Montserrat Thin Italic'), url('fonts/Montserrat-ThinItalic.woff') format('woff'); -} +/*@font-face {*/ +/* font-family: 'Montserrat Thin Italic';*/ +/* font-style: normal;*/ +/* font-weight: normal;*/ +/* src: local('Montserrat Thin Italic'), url('fonts/montserrat-thinitalic.woff') format('woff');*/ +/*}*/ -@font-face { - font-family: 'Montserrat ExtraLight'; - font-style: normal; - font-weight: normal; - src: local('Montserrat ExtraLight'), url('fonts/Montserrat-ExtraLight.woff') format('woff'); -} +/*@font-face {*/ +/* font-family: 'Montserrat ExtraLight';*/ +/* font-style: normal;*/ +/* font-weight: normal;*/ +/* src: local('Montserrat ExtraLight'), url('fonts/Montserrat-ExtraLight.woff') format('woff');*/ +/*}*/ -@font-face { - font-family: 'Montserrat ExtraLight Italic'; - font-style: normal; - font-weight: normal; - src: local('Montserrat ExtraLight Italic'), url('fonts/Montserrat-ExtraLightItalic.woff') format('woff'); -} +/*@font-face {*/ +/* font-family: 'Montserrat ExtraLight Italic';*/ +/* font-style: normal;*/ +/* font-weight: normal;*/ +/* src: local('Montserrat ExtraLight Italic'), url('fonts/Montserrat-ExtraLightItalic.woff') format('woff');*/ +/*}*/ -@font-face { - font-family: 'Montserrat Light'; - font-style: normal; - font-weight: normal; - src: local('Montserrat Light'), url('fonts/Montserrat-Light.woff') format('woff'); -} +/*@font-face {*/ +/* font-family: 'Montserrat Light';*/ +/* font-style: normal;*/ +/* font-weight: normal;*/ +/* src: local('Montserrat Light'), url('fonts/Montserrat-Light.woff') format('woff');*/ +/*}*/ -@font-face { - font-family: 'Montserrat Light Italic'; - font-style: normal; - font-weight: normal; - src: local('Montserrat Light Italic'), url('fonts/Montserrat-LightItalic.woff') format('woff'); -} +/*@font-face {*/ +/* font-family: 'Montserrat Light Italic';*/ +/* font-style: normal;*/ +/* font-weight: normal;*/ +/* src: local('Montserrat Light Italic'), url('fonts/Montserrat-LightItalic.woff') format('woff');*/ +/*}*/ -@font-face { - font-family: 'Montserrat Medium'; - font-style: normal; - font-weight: normal; - src: local('Montserrat Medium'), url('fonts/Montserrat-Medium.woff') format('woff'); -} +/*@font-face {*/ +/* font-family: 'Montserrat Medium';*/ +/* font-style: normal;*/ +/* font-weight: normal;*/ +/* src: local('Montserrat Medium'), url('fonts/Montserrat-Medium.woff') format('woff');*/ +/*}*/ -@font-face { - font-family: 'Montserrat Medium Italic'; - font-style: normal; - font-weight: normal; - src: local('Montserrat Medium Italic'), url('fonts/Montserrat-MediumItalic.woff') format('woff'); -} +/*@font-face {*/ +/* font-family: 'Montserrat Medium Italic';*/ +/* font-style: normal;*/ +/* font-weight: normal;*/ +/* src: local('Montserrat Medium Italic'), url('fonts/Montserrat-MediumItalic.woff') format('woff');*/ +/*}*/ -@font-face { - font-family: 'Montserrat SemiBold'; - font-style: normal; - font-weight: normal; - src: local('Montserrat SemiBold'), url('fonts/Montserrat-SemiBold.woff') format('woff'); -} +/*@font-face {*/ +/* font-family: 'Montserrat SemiBold';*/ +/* font-style: normal;*/ +/* font-weight: normal;*/ +/* src: local('Montserrat SemiBold'), url('fonts/Montserrat-SemiBold.woff') format('woff');*/ +/*}*/ -@font-face { - font-family: 'Montserrat SemiBold Italic'; - font-style: normal; - font-weight: normal; - src: local('Montserrat SemiBold Italic'), url('fonts/Montserrat-SemiBoldItalic.woff') format('woff'); -} +/*@font-face {*/ +/* font-family: 'Montserrat SemiBold Italic';*/ +/* font-style: normal;*/ +/* font-weight: normal;*/ +/* src: local('Montserrat SemiBold Italic'), url('fonts/Montserrat-SemiBoldItalic.woff') format('woff');*/ +/*}*/ -@font-face { - font-family: 'Montserrat Bold'; - font-style: normal; - font-weight: normal; - src: local('Montserrat Bold'), url('fonts/Montserrat-Bold.woff') format('woff'); -} +/*@font-face {*/ +/* font-family: 'Montserrat Bold';*/ +/* font-style: normal;*/ +/* font-weight: normal;*/ +/* src: local('Montserrat Bold'), url('fonts/Montserrat-Bold.woff') format('woff');*/ +/*}*/ -@font-face { - font-family: 'Montserrat Bold Italic'; - font-style: normal; - font-weight: normal; - src: local('Montserrat Bold Italic'), url('fonts/Montserrat-BoldItalic.woff') format('woff'); -} +/*@font-face {*/ +/* font-family: 'Montserrat Bold Italic';*/ +/* font-style: normal;*/ +/* font-weight: normal;*/ +/* src: local('Montserrat Bold Italic'), url('fonts/Montserrat-BoldItalic.woff') format('woff');*/ +/*}*/ -@font-face { - font-family: 'Montserrat ExtraBold'; - font-style: normal; - font-weight: normal; - src: local('Montserrat ExtraBold'), url('fonts/Montserrat-ExtraBold.woff') format('woff'); -} +/*@font-face {*/ +/* font-family: 'Montserrat ExtraBold';*/ +/* font-style: normal;*/ +/* font-weight: normal;*/ +/* src: local('Montserrat ExtraBold'), url('fonts/Montserrat-ExtraBold.woff') format('woff');*/ +/*}*/ -@font-face { - font-family: 'Montserrat ExtraBold Italic'; - font-style: normal; - font-weight: normal; - src: local('Montserrat ExtraBold Italic'), url('fonts/Montserrat-ExtraBoldItalic.woff') format('woff'); -} +/*@font-face {*/ +/* font-family: 'Montserrat ExtraBold Italic';*/ +/* font-style: normal;*/ +/* font-weight: normal;*/ +/* src: local('Montserrat ExtraBold Italic'), url('fonts/Montserrat-ExtraBoldItalic.woff') format('woff');*/ +/*}*/ -@font-face { - font-family: 'Montserrat Black'; - font-style: normal; - font-weight: normal; - src: local('Montserrat Black'), url('fonts/Montserrat-Black.woff') format('woff'); -} +/*@font-face {*/ +/* font-family: 'Montserrat Black';*/ +/* font-style: normal;*/ +/* font-weight: normal;*/ +/* src: local('Montserrat Black'), url('fonts/Montserrat-Black.woff') format('woff');*/ +/*}*/ -@font-face { - font-family: 'Montserrat Black Italic'; - font-style: normal; - font-weight: normal; - src: local('Montserrat Black Italic'), url('fonts/Montserrat-BlackItalic.woff') format('woff'); -} \ No newline at end of file +/*@font-face {*/ +/* font-family: 'Montserrat Black Italic';*/ +/* font-style: normal;*/ +/* font-weight: normal;*/ +/* src: local('Montserrat Black Italic'), url('fonts/Montserrat-BlackItalic.woff') format('woff');*/ +/*}*/ \ No newline at end of file diff --git a/frontend/public/assets/fonts/montserrat-regular.woff b/frontend/public/assets/fonts/montserrat-regular.woff new file mode 100644 index 000000000..2a990e5d8 Binary files /dev/null and b/frontend/public/assets/fonts/montserrat-regular.woff differ diff --git a/frontend/public/assets/i18n/en.json b/frontend/public/assets/i18n/en.json index 84cf146f4..d0e15d7e8 100644 --- a/frontend/public/assets/i18n/en.json +++ b/frontend/public/assets/i18n/en.json @@ -136,6 +136,7 @@ "NOTES": "Notes", "DASHBOARD": "Dashboard", "STATS": "Stats", + "STREAMS": "Streams", "SETTINGS": "Settings", "WELCOME": "Welcome", "RELEASES": "Releases", @@ -283,7 +284,12 @@ "ADD_EXTENSION": "Add Extension", "ADD_DEVICE": "Add Device", "EXTENDED_ATTRIBUTES": "Extended Attributes", - "USE_MEMORY_CACHE": "Mem-Cache" + "USE_MEMORY_CACHE": "Mem-Cache", + "KICK": "Kick", + "CHANNEL": "Channel", + "GROUP": "Group", + "CLIENT_IP": "Client IP", + "STREAM_ID": "Stream Id" }, "TITLE": { "USER_BOUQUET_EDITOR": "User group editor" diff --git a/frontend/public/assets/icons.json b/frontend/public/assets/icons.json index 5a6178e73..83f7d9da7 100644 --- a/frontend/public/assets/icons.json +++ b/frontend/public/assets/icons.json @@ -503,5 +503,13 @@ { "keys": ["Epg"], "path": "M21 3H3c-1.1 0-2 .9-2 2v12c0 1.1.9 2 2 2h5v2h8v-2h5c1.1 0 1.99-.9 1.99-2L23 5c0-1.1-.9-2-2-2m0 14H3V5h18zm-2-9H8v2h11zm0 4H8v2h11zM7 8H5v2h2zm0 4H5v2h2z" + }, + { + "keys": ["Streams"], + "path": "m 17.440433,1.9998779 c -0.518792,0 -0.868751,0.3408817 -0.868751,0.8464598 0,0.3202564 0.155237,0.5904331 0.7372,1.2820925 3.940613,4.6833925 3.827549,11.3636248 -0.268776,15.8946368 -0.775588,0.857895 -0.789451,1.433059 -0.0454,1.840715 0.13286,0.07279 0.284785,0.134108 0.337904,0.135909 0.616589,0.02349 2.293279,-2.001409 3.252654,-3.928443 C 23.079301,13.061613 22.210244,6.6790696 18.495417,2.7176634 17.94865,2.1346036 17.751133,1.9998779 17.440433,1.9998779 Z m -3.734491,3.9051878 c -0.198863,0.00493 -0.385897,0.1075021 -0.581918,0.3038576 -0.428044,0.4287782 -0.375106,0.807358 0.207386,1.480013 1.198656,1.3841968 1.669919,2.5734407 1.669919,4.2121457 0,1.638701 -0.471263,2.827437 -1.669919,4.21163 -0.569993,0.658217 -0.65316,1.144066 -0.262586,1.535306 0.14662,0.146869 0.402874,0.241846 0.651563,0.241846 0.345263,0 0.514191,-0.113051 1.074072,-0.718819 C 17.589057,14.147402 17.555625,9.4820493 14.718108,6.5225993 14.312382,6.0994398 13.996588,5.8978678 13.705942,5.9050657 Z M 3.3356259,6.3975424 C 3.1030533,6.3976138 2.8707078,6.4500288 2.6685868,6.5546387 2.256259,6.7674812 2.0000001,7.1652814 2,7.5923012 v 8.8118648 c 0,0.427212 0.2560217,0.824932 0.6660073,1.03663 0.2027648,0.105864 0.4346085,0.161747 0.6701344,0.161747 0.2348851,0 0.4661626,-0.05562 0.6680709,-0.160714 L 12.518375,13.03538 c 0.411899,-0.213414 0.667555,-0.610701 0.667555,-1.037146 -2.1e-4,-0.425877 -0.255747,-0.823107 -0.66807,-1.037663 L 4.0036967,6.5546387 C 3.8008304,6.4497431 3.5681984,6.3974707 3.3356259,6.3975424 Z m 0,1.1410156 c 0.010777,0 0.021367,0.00247 0.030437,0.00723 l 8.5126151,4.404904 c 0.01902,0.01 0.03095,0.02809 0.03095,0.04754 -2.13e-4,0.01971 -0.0114,0.03729 -0.02992,0.04703 v -5.17e-4 l -8.5146787,4.406449 c -0.017543,0.0089 -0.038274,0.0099 -0.059327,-0.001 -0.017969,-0.0093 -0.029406,-0.0269 -0.029405,-0.04599 V 7.5917845 c 0,-0.01908 0.011131,-0.036264 0.029921,-0.045992 0.00886,-0.00477 0.019069,-0.00723 0.029405,-0.00723 z" + }, + { + "keys": ["Disconnect"], + "path": "M3.42 2.36L2.01 3.78 4.1 5.87C2.79 7.57 2 9.69 2 12c0 3.7 2.01 6.92 4.99 8.65l1-1.73C5.61 17.53 4 14.96 4 12c0-1.76.57-3.38 1.53-4.69l1.43 1.44C6.36 9.68 6 10.8 6 12c0 2.22 1.21 4.15 3 5.19l1-1.74c-1.19-.7-2-1.97-2-3.45 0-.65.17-1.25.44-1.79l1.58 1.58L10 12c0 1.1.9 2 2 2l.21-.02 7.52 7.52 1.41-1.41L3.42 2.36zm14.29 11.46c.18-.57.29-1.19.29-1.82 0-3.31-2.69-6-6-6-.63 0-1.25.11-1.82.29l1.72 1.72c.03 0 .06-.01.1-.01 2.21 0 4 1.79 4 4 0 .04-.01.07-.01.11l1.72 1.71zM12 4c4.42 0 8 3.58 8 8 0 1.2-.29 2.32-.77 3.35l1.49 1.49C21.53 15.4 22 13.76 22 12c0-5.52-4.48-10-10-10-1.76 0-3.4.48-4.84 1.28l1.48 1.48C9.66 4.28 10.8 4 12 4z" } ] \ No newline at end of file diff --git a/frontend/scss/app/_component.scss b/frontend/scss/app/_component.scss index 1ed647ebc..e60e68382 100644 --- a/frontend/scss/app/_component.scss +++ b/frontend/scss/app/_component.scss @@ -17,6 +17,7 @@ @forward "components/dashboard/status_card"; @forward "components/dashboard/dashboard_view"; @forward "components/dashboard/stats_view"; +@forward "components/dashboard/streams_view"; @forward "components/list_view"; @forward "components/reveal_content"; @forward "components/hide_content"; diff --git a/frontend/scss/app/components/_no_content.scss b/frontend/scss/app/components/_no_content.scss index 16f383908..3095a99c4 100644 --- a/frontend/scss/app/components/_no_content.scss +++ b/frontend/scss/app/components/_no_content.scss @@ -5,6 +5,7 @@ align-items: center; padding: var(--padding-default); width: 100%; + box-sizing: border-box; &__indicator { display: flex; diff --git a/frontend/scss/app/components/dashboard/_streams_view.scss b/frontend/scss/app/components/dashboard/_streams_view.scss new file mode 100644 index 000000000..8ad4295eb --- /dev/null +++ b/frontend/scss/app/components/dashboard/_streams_view.scss @@ -0,0 +1,21 @@ +.tp__streams { + display: flex; + flex-flow: column; + flex: 1 1 auto; + box-sizing: border-box; + width: 100%; + max-width: var(--max-view-width); + overflow: hidden; + + &__header { + display: flex; + flex-flow: row nowrap; + } + + &__body { + display: flex; + flex-flow: column; + gap: var(--gap-larger); + overflow: auto; + } +} \ No newline at end of file diff --git a/frontend/scss/main.scss b/frontend/scss/main.scss index 981b2fc45..60b709081 100644 --- a/frontend/scss/main.scss +++ b/frontend/scss/main.scss @@ -8,7 +8,7 @@ html { body { margin: 0; padding: 0; - font-family: "Montesserat", sans-serif; + font-family: "Montserrat", sans-serif; font-size: 14px; background-color: var(--background-color); height: 100%; diff --git a/frontend/src/app/components/dashboard/mod.rs b/frontend/src/app/components/dashboard/mod.rs index 9caac5d9f..71c9c3a29 100644 --- a/frontend/src/app/components/dashboard/mod.rs +++ b/frontend/src/app/components/dashboard/mod.rs @@ -10,6 +10,9 @@ mod dashboard_view; mod stats_view; mod playlist_progress_status_card; +mod streams_view; +mod streams_table; + pub use self::action_card::*; pub use self::status_card::*; pub use self::user_action_card::*; @@ -21,3 +24,5 @@ pub use self::github_action_card::*; pub use self::dashboard_view::*; pub use self::stats_view::*; pub use self::playlist_progress_status_card::*; +pub use self::streams_view::*; +pub use self::streams_table::*; diff --git a/frontend/src/app/components/dashboard/streams_table.rs b/frontend/src/app/components/dashboard/streams_table.rs new file mode 100644 index 000000000..55e700c69 --- /dev/null +++ b/frontend/src/app/components/dashboard/streams_table.rs @@ -0,0 +1,191 @@ +use crate::app::components::menu_item::MenuItem; +use crate::app::components::popup_menu::PopupMenu; +use crate::app::components::{AppIcon, Table, TableDefinition}; +use crate::hooks::use_service_context; +use shared::error::{create_tuliprox_error_result, TuliproxError, TuliproxErrorKind}; +use shared::model::{SortOrder, StreamInfo}; +use std::fmt::Display; +use std::rc::Rc; +use std::str::FromStr; +use yew::prelude::*; +use yew_i18n::use_translation; + +const HEADERS: [&str; 7] = [ + "LABEL.EMPTY", + "LABEL.USERNAME", + "LABEL.STREAM_ID", + "LABEL.CHANNEL", + "LABEL.GROUP", + "LABEL.CLIENT_IP", + "LABEL.PROVIDER" +]; + +#[derive(Properties, PartialEq, Clone)] +pub struct StreamsTableProps { + pub streams: Option>>, +} + +#[function_component] +pub fn StreamsTable(props: &StreamsTableProps) -> Html { + let translate = use_translation(); + let services = use_service_context(); + let popup_anchor_ref = use_state(|| None::); + let popup_is_open = use_state(|| false); + let selected_dto = use_state(|| None::>); + + let handle_popup_close = { + let set_is_open = popup_is_open.clone(); + Callback::from(move |()| { + set_is_open.set(false); + }) + }; + + let handle_popup_onclick = { + let set_selected_dto = selected_dto.clone(); + let set_anchor_ref = popup_anchor_ref.clone(); + let set_is_open = popup_is_open.clone(); + Callback::from(move |(dto, event): (Rc, MouseEvent)| { + if let Some(streams) = event.target_dyn_into::() { + set_selected_dto.set(Some(dto.clone())); + set_anchor_ref.set(Some(streams)); + set_is_open.set(true); + } + }) + }; + + let render_header_cell = { + let translator = translate.clone(); + Callback::::from(move |col| { + html! { + { + if col < HEADERS.len() { + translator.t(HEADERS[col]) + } else { + String::new() + } + } + } + }) + }; + + let render_data_cell = { + let popup_onclick = handle_popup_onclick.clone(); + Callback::<(usize, usize, Rc), Html>::from( + move |(row, col, dto): (usize, usize, Rc)| { + match col { + 0 => { + let popup_onclick = popup_onclick.clone(); + html! { + + } + } + 1 => html! {dto.username.as_str()}, + 2 => html! { <> + { dto.channel.virtual_id.to_string() } + {" ("} + { dto.channel.provider_id.to_string() } + {")"} + }, + 3 => html! {dto.channel.title.as_str()}, + 4 => html! {dto.channel.group.as_str()}, + 5 => html! {dto.addr.as_str()}, + 6 => html! {dto.provider.as_str()}, + _ => html! {""}, + } + }) + }; + + let is_sortable = Callback::::from(move |_col| { + false + }); + + let on_sort = Callback::, ()>::from(move |_args| { + }); + + let table_definition = { + // first register for config update + let render_header_cell_cb = render_header_cell.clone(); + let render_data_cell_cb = render_data_cell.clone(); + let is_sortable = is_sortable.clone(); + let on_sort = on_sort.clone(); + let num_cols = HEADERS.len(); + use_memo(props.streams.clone(), move |streams| { + streams.as_ref().map(|list| + Rc::new(TableDefinition:: { + items: if list.is_empty() {None} else {Some(Rc::new(list.clone()))}, + num_cols, + is_sortable, + on_sort, + render_header_cell: render_header_cell_cb, + render_data_cell: render_data_cell_cb, + })) + }) + }; + + + let handle_menu_click = { + let popup_is_open_state = popup_is_open.clone(); + //let translate = translate.clone(); + let services_ctx = services.clone(); + //let selected_dto = selected_dto.clone(); + Callback::from(move |(name, _): (String, _)| { + if let Ok(action) = StreamsTableAction::from_str(&name) { + match action { + StreamsTableAction::Kick => { + // TODO implement connection kick + services_ctx.toastr.error("Not implemented") + } + } + } + popup_is_open_state.set(false); + }) + }; + + html! { +
+ { + if let Some(definition) = table_definition.as_ref() { + html! { + <> + definition={definition.clone()} /> + + + + + } + } else { + html! {} + } + } +
+ } +} + +#[derive(Debug, Clone, Eq, PartialEq)] +enum StreamsTableAction { + Kick, +} + +impl Display for StreamsTableAction { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + write!(f, "{}", match self { + Self::Kick => "kick", + }) + } +} + +impl FromStr for StreamsTableAction { + type Err = TuliproxError; + + fn from_str(s: &str) -> Result { + if s.eq("kick") { + Ok(Self::Kick) + } else { + create_tuliprox_error_result!(TuliproxErrorKind::Info, "Unknown Stream Action: {}", s) + } + } +} \ No newline at end of file diff --git a/frontend/src/app/components/dashboard/streams_view.rs b/frontend/src/app/components/dashboard/streams_view.rs new file mode 100644 index 000000000..49c9c6e92 --- /dev/null +++ b/frontend/src/app/components/dashboard/streams_view.rs @@ -0,0 +1,22 @@ +use std::rc::Rc; +use yew::prelude::*; +use yew_i18n::use_translation; +use crate::app::components::StreamsTable; +use crate::app::StatusContext; + +#[function_component] +pub fn StreamsView() -> Html { + let translate = use_translation(); + let status_ctx = use_context::().expect("Status context not found"); + + html! { +
+
+

{ translate.t("LABEL.STREAMS")}

+
+
+ +
+
+ } +} \ No newline at end of file diff --git a/frontend/src/app/components/home.rs b/frontend/src/app/components/home.rs index acd16429e..45e6b19c2 100644 --- a/frontend/src/app/components/home.rs +++ b/frontend/src/app/components/home.rs @@ -1,4 +1,4 @@ -use crate::app::components::{AppIcon, DashboardView, EpgView, IconButton, InputRow, Panel, PlaylistEditorView, PlaylistExplorerView, PlaylistUpdateView, Sidebar, StatsView, ToastrView, UserlistView, WebsocketStatus}; +use crate::app::components::{AppIcon, DashboardView, EpgView, IconButton, InputRow, Panel, PlaylistEditorView, PlaylistExplorerView, PlaylistUpdateView, Sidebar, StatsView, StreamsView, ToastrView, UserlistView, WebsocketStatus}; use crate::app::context::{ConfigContext, PlaylistContext, StatusContext}; use crate::hooks::{use_server_status, use_service_context}; use crate::model::{EventMessage, ViewType}; @@ -135,6 +135,7 @@ pub fn Home() -> Html { sources: sources.clone(), }; + //
html! { @@ -171,6 +172,9 @@ pub fn Home() -> Html { + + + diff --git a/frontend/src/app/components/playlist/target_table.rs b/frontend/src/app/components/playlist/target_table.rs index b4ac9398e..01abb0d04 100644 --- a/frontend/src/app/components/playlist/target_table.rs +++ b/frontend/src/app/components/playlist/target_table.rs @@ -151,10 +151,10 @@ pub fn TargetTable(props: &TargetTableProps) -> Html { let services_ctx = services.clone(); let selected_dto = selected_dto.clone(); Callback::from(move |(name, _): (String, _)| { - if let Ok(action) = TableAction::from_str(&name) { + if let Ok(action) = TargetTableAction::from_str(&name) { match action { - TableAction::Edit => {} - TableAction::Refresh => { + TargetTableAction::Edit => {} + TargetTableAction::Refresh => { let translate = translate.clone(); let services_ctx = services_ctx.clone(); let dto_name = selected_dto.as_ref().map_or_else(String::new, |d| d.name.to_string()); @@ -166,7 +166,7 @@ pub fn TargetTable(props: &TargetTableProps) -> Html { } }); } - TableAction::Delete => { + TargetTableAction::Delete => { let confirm = confirm.clone(); let translator = translate.clone(); spawn_local(async move { @@ -190,10 +190,10 @@ pub fn TargetTable(props: &TargetTableProps) -> Html { <> definition={definition.clone()} /> - - + +
- +
} @@ -207,13 +207,13 @@ pub fn TargetTable(props: &TargetTableProps) -> Html { #[derive(Debug, Clone, Eq, PartialEq)] -enum TableAction { +enum TargetTableAction { Edit, Refresh, Delete, } -impl Display for TableAction { +impl Display for TargetTableAction { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { write!(f, "{}", match self { Self::Edit => "edit", @@ -223,7 +223,7 @@ impl Display for TableAction { } } -impl FromStr for TableAction { +impl FromStr for TargetTableAction { type Err = TuliproxError; fn from_str(s: &str) -> Result { @@ -234,7 +234,7 @@ impl FromStr for TableAction { } else if s.eq("delete") { Ok(Self::Delete) } else { - create_tuliprox_error_result!(TuliproxErrorKind::Info, "Unknown InputType: {}", s) + create_tuliprox_error_result!(TuliproxErrorKind::Info, "Unknown Target Action: {}", s) } } } \ No newline at end of file diff --git a/frontend/src/app/components/sidebar.rs b/frontend/src/app/components/sidebar.rs index 1cdd70c20..7354f3f27 100644 --- a/frontend/src/app/components/sidebar.rs +++ b/frontend/src/app/components/sidebar.rs @@ -135,6 +135,7 @@ pub fn Sidebar(props: &SidebarProps) -> Html {
+ @@ -154,6 +155,7 @@ pub fn Sidebar(props: &SidebarProps) -> Html {
+ diff --git a/frontend/src/hooks/use_server_status.rs b/frontend/src/hooks/use_server_status.rs index 78695d72b..bf5568f6a 100644 --- a/frontend/src/hooks/use_server_status.rs +++ b/frontend/src/hooks/use_server_status.rs @@ -3,7 +3,7 @@ use std::rc::Rc; use std::collections::BTreeMap; use gloo_timers::callback::Interval; use yew::prelude::*; -use shared::model::StatusCheck; +use shared::model::{ActiveUserConnectionChange, StatusCheck}; use crate::hooks::use_service_context; use yew::platform::spawn_local; use crate::model::EventMessage; @@ -27,7 +27,7 @@ pub fn use_server_status( *status_holder_signal.borrow_mut() = Some(Rc::clone(&server_status)); status_signal.set(Some(server_status)); } - EventMessage::ActiveUser(user_count, connections) => { + EventMessage::ActiveUser(event) => { let mut server_status = { if let Some(old_status) = status_holder_signal.borrow().as_ref() { (**old_status).clone() @@ -35,8 +35,20 @@ pub fn use_server_status( StatusCheck::default() } }; - server_status.active_users = user_count; - server_status.active_user_connections = connections; + + match event { + ActiveUserConnectionChange::Connected(stream_info) => { + server_status.active_user_streams.push(stream_info); + } + ActiveUserConnectionChange::Disconnected(addr) => { + server_status.active_user_streams.retain(|stream_info| stream_info.addr != addr); + } + ActiveUserConnectionChange::Connections(user_count, connections) => { + server_status.active_users = user_count; + server_status.active_user_connections = connections; + } + } + let new_status = Rc::new(server_status); *status_holder_signal.borrow_mut() = Some(Rc::clone(&new_status)); status_signal.set(Some(new_status)); diff --git a/frontend/src/hooks/use_service_context.rs b/frontend/src/hooks/use_service_context.rs index 0ea6e633b..d0463f154 100644 --- a/frontend/src/hooks/use_service_context.rs +++ b/frontend/src/hooks/use_service_context.rs @@ -1,14 +1,14 @@ use std::rc::Rc; use yew::prelude::*; use crate::model::WebConfig; -use crate::services::{AuthService, ConfigService, EventService, PlaylistService, StatusService, ToastrService, - UserService, WebSocketService}; +use crate::services::{AuthService, ConfigService, EventService, PlaylistService, StatusService, StreamsService, ToastrService, UserService, WebSocketService}; pub struct Services { pub auth: Rc, pub config: Rc, pub user: Rc, pub status: Rc, + pub streams: Rc, pub event: Rc, pub playlist: Rc, pub toastr: Rc, @@ -21,6 +21,7 @@ impl Services { let config = Rc::new(ConfigService::new(web_config, Rc::clone(&event))); let auth = Rc::new(AuthService::new()); let status = Rc::new(StatusService::new()); + let streams = Rc::new(StreamsService::new()); let playlist = Rc::new(PlaylistService::new()); let toastr = Rc::new(ToastrService::new()); let user = Rc::new(UserService::new(Rc::clone(&event))); @@ -29,6 +30,7 @@ impl Services { auth, config, status, + streams, event, playlist, user, diff --git a/frontend/src/model/event_message.rs b/frontend/src/model/event_message.rs index 7061a6be7..7bbc3c45a 100644 --- a/frontend/src/model/event_message.rs +++ b/frontend/src/model/event_message.rs @@ -1,5 +1,5 @@ use std::rc::Rc; -use shared::model::{ConfigType, PlaylistUpdateState, StatusCheck}; +use shared::model::{ActiveUserConnectionChange, ConfigType, PlaylistUpdateState, StatusCheck}; use crate::model::BusyStatus; #[derive(Debug, Clone, PartialEq)] @@ -7,7 +7,7 @@ pub enum EventMessage { Unauthorized, ServerError(String), ServerStatus(Rc), - ActiveUser(usize, usize), + ActiveUser(ActiveUserConnectionChange), ActiveProvider(String, usize), ConfigChange(ConfigType), Busy(BusyStatus), diff --git a/frontend/src/model/view_type.rs b/frontend/src/model/view_type.rs index 3e47a1553..e61044d87 100644 --- a/frontend/src/model/view_type.rs +++ b/frontend/src/model/view_type.rs @@ -4,6 +4,7 @@ use shared::error::{info_err, TuliproxError}; const DASHBOARD: &str = "dashboard"; const STATS: &str = "stats"; +const STREAMS: &str = "streams"; const USERS: &str = "users"; const CONFIG: &str = "config"; const PLAYLIST_UPDATE: &str = "playlist_update"; @@ -16,6 +17,7 @@ const PLAYLIST_EPG: &str = "playlist_epg"; pub enum ViewType { Dashboard, Stats, + Streams, Users, Config, PlaylistUpdate, @@ -31,6 +33,7 @@ impl FromStr for ViewType { match s.to_lowercase().as_str() { DASHBOARD => Ok(ViewType::Dashboard), STATS => Ok(ViewType::Stats), + STREAMS => Ok(ViewType::Streams), USERS => Ok(ViewType::Users), CONFIG => Ok(ViewType::Config), PLAYLIST_UPDATE => Ok(ViewType::PlaylistUpdate), @@ -47,6 +50,7 @@ impl fmt::Display for ViewType { let s = match self { ViewType::Dashboard => DASHBOARD, ViewType::Stats => STATS, + ViewType::Streams => STREAMS, ViewType::Users => USERS, ViewType::Config => CONFIG, ViewType::PlaylistUpdate => PLAYLIST_UPDATE, diff --git a/frontend/src/services/mod.rs b/frontend/src/services/mod.rs index c0c3dfa25..25095df33 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 streams_service; pub use self::auth_service::*; pub use self::config_service::*; @@ -18,4 +19,5 @@ pub use self::playlist_service::*; pub use self::websocket_service::*; pub use self::toastr_service::*; pub use self::event_service::*; -pub use self::user_service::*; \ No newline at end of file +pub use self::user_service::*; +pub use self::streams_service::*; \ No newline at end of file diff --git a/frontend/src/services/streams_service.rs b/frontend/src/services/streams_service.rs new file mode 100644 index 000000000..dd7d29b9c --- /dev/null +++ b/frontend/src/services/streams_service.rs @@ -0,0 +1,27 @@ +use std::rc::Rc; +use crate::services::{get_base_href, request_get}; +use shared::model::StreamInfo; +use shared::utils::concat_path_leading_slash; + +pub struct StreamsService { + streams_path: String, +} + +impl Default for StreamsService { + fn default() -> Self { + Self::new() + } +} + +impl StreamsService { + pub fn new() -> Self { + let base_href = get_base_href(); + Self { + streams_path: concat_path_leading_slash(&base_href, "api/v1/streams"), + } + } + + pub async fn get_streams_info(&self) -> Result>>, crate::error::Error> { + request_get::>>(&self.streams_path, None, None).await + } +} \ No newline at end of file diff --git a/frontend/src/services/websocket_service.rs b/frontend/src/services/websocket_service.rs index b685d41a3..27158588c 100644 --- a/frontend/src/services/websocket_service.rs +++ b/frontend/src/services/websocket_service.rs @@ -72,8 +72,8 @@ impl WebSocketService { ProtocolMessage::Error(err) => { error!("{err}"); }, - ProtocolMessage::ActiveUserResponse(user_count, connections) => { - event_service.broadcast(EventMessage::ActiveUser(user_count, connections)); + ProtocolMessage::ActiveUserResponse(event) => { + event_service.broadcast(EventMessage::ActiveUser(event)); }, ProtocolMessage::ActiveProviderResponse(user_count, connections) => { event_service.broadcast(EventMessage::ActiveProvider(user_count, connections)); diff --git a/shared/src/model/active_user_connection_change.rs b/shared/src/model/active_user_connection_change.rs new file mode 100644 index 000000000..b67451cd1 --- /dev/null +++ b/shared/src/model/active_user_connection_change.rs @@ -0,0 +1,8 @@ +use crate::model::StreamInfo; + +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, PartialEq)] +pub enum ActiveUserConnectionChange { + Connected(StreamInfo), + Disconnected(String), // addr + Connections(usize, usize) // user_count, connection_count +} diff --git a/shared/src/model/mod.rs b/shared/src/model/mod.rs index ab5535de9..cc8dc0301 100644 --- a/shared/src/model/mod.rs +++ b/shared/src/model/mod.rs @@ -18,6 +18,8 @@ mod search_request; mod webplayer_url_request; mod epg; mod epg_request; +mod stream_info; +mod active_user_connection_change; pub use self::cluster_flags::*; pub use self::playlist::*; @@ -38,3 +40,5 @@ pub use self::search_request::*; pub use self::webplayer_url_request::*; pub use self::epg::*; pub use self::epg_request::*; +pub use self::stream_info::*; +pub use self::active_user_connection_change::*; diff --git a/shared/src/model/playlist.rs b/shared/src/model/playlist.rs index d8d22ba4d..032cd5ee0 100644 --- a/shared/src/model/playlist.rs +++ b/shared/src/model/playlist.rs @@ -158,17 +158,11 @@ pub struct PlaylistItemHeader { #[serde(default)] pub category_id: u32, pub input_name: String, - #[serde(default)] - pub copy: bool, // not original, a copy } impl PlaylistItemHeader { pub fn gen_uuid(&mut self) { - self.uuid = if self.copy { - generate_playlist_uuid(&format!("copy-{}", self.input_name), &self.id, self.item_type, &format!("copy-{}", self.url)) - } else { - generate_playlist_uuid(&self.input_name, &self.id, self.item_type, &self.url) - }; + self.uuid = generate_playlist_uuid(&self.input_name, &self.id, self.item_type, &self.url); } pub const fn get_uuid(&self) -> &UUIDType { &self.uuid @@ -299,8 +293,6 @@ pub struct M3uPlaylistItem { pub epg_channel_id: Option, pub input_name: String, pub item_type: PlaylistItemType, - #[serde(default)] - pub copy: bool, #[serde(skip)] pub t_stream_url: String, #[serde(skip)] @@ -433,8 +425,6 @@ pub struct XtreamPlaylistItem { pub category_id: u32, pub input_name: String, pub channel_no: u32, - #[serde(default)] - pub copy: bool, } impl XtreamPlaylistItem { @@ -595,7 +585,6 @@ impl PlaylistItem { epg_channel_id: header.epg_channel_id.clone(), input_name: header.input_name.to_string(), item_type: header.item_type, - copy: header.copy, t_stream_url: header.url.to_string(), t_resource_url: None, } @@ -655,7 +644,6 @@ impl PlaylistItem { category_id: header.category_id, input_name: header.input_name.to_string(), channel_no: header.chno.parse::().unwrap_or(0), - copy: header.copy, } } @@ -771,4 +759,4 @@ impl PlaylistGroup { { self.channels.iter().filter(|&c| filter(c)).count() } -} \ No newline at end of file +} diff --git a/shared/src/model/status_check.rs b/shared/src/model/status_check.rs index 2c2b36784..bfa4f7b45 100644 --- a/shared/src/model/status_check.rs +++ b/shared/src/model/status_check.rs @@ -1,5 +1,6 @@ use std::collections::BTreeMap; use serde::{Deserialize, Serialize}; +use crate::model::StreamInfo; #[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] pub struct StatusCheck { @@ -13,6 +14,7 @@ pub struct StatusCheck { pub cache: Option, pub active_users: usize, pub active_user_connections: usize, + pub active_user_streams: Vec, #[serde(skip_serializing_if = "Option::is_none")] pub active_provider_connections: Option>, } @@ -29,6 +31,7 @@ impl Default for StatusCheck { active_users: 0, active_user_connections: 0, active_provider_connections: None, + active_user_streams: Vec::new(), } } } \ No newline at end of file diff --git a/shared/src/model/stream_info.rs b/shared/src/model/stream_info.rs new file mode 100644 index 000000000..6cbb39d3f --- /dev/null +++ b/shared/src/model/stream_info.rs @@ -0,0 +1,59 @@ +use serde::{Deserialize, Serialize}; +use crate::model::{M3uPlaylistItem, PlaylistEntry, PlaylistItemType, XtreamCluster, XtreamPlaylistItem}; + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] +pub struct StreamChannel { + pub virtual_id: u32, + pub provider_id: u32, + pub item_type: PlaylistItemType, + pub cluster: XtreamCluster, + pub group: String, + pub title: String, + pub url: String, +} +impl XtreamPlaylistItem { + pub fn to_stream_channel(&self) -> StreamChannel { + StreamChannel { + virtual_id: self.virtual_id, + provider_id: self.provider_id, + item_type: self.item_type, + cluster: self.xtream_cluster, + group: self.group.clone(), + title: self.title.clone(), + url: self.url.clone(), + } + } +} + +impl M3uPlaylistItem { + pub fn to_stream_channel(&self) -> StreamChannel { + StreamChannel { + virtual_id: self.virtual_id, + provider_id: self.get_provider_id().unwrap_or_default(), + item_type: self.item_type, + cluster: XtreamCluster::try_from(self.item_type).unwrap_or(XtreamCluster::Live), + group: self.group.clone(), + title: self.title.clone(), + url: self.url.clone(), + } + } +} + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] +pub struct StreamInfo { + pub username: String, + pub channel: StreamChannel, + pub provider: String, + pub addr: String, +} + +impl StreamInfo { + pub fn new(username: &str, addr: &str, provider: &str, stream_channel: StreamChannel) -> Self { + Self { + username: username.to_string(), + channel: stream_channel, + provider: provider.to_string(), + addr: addr.to_string(), + } + } +} \ No newline at end of file diff --git a/shared/src/model/web_socket.rs b/shared/src/model/web_socket.rs index 0cf677c5a..6af24ac92 100644 --- a/shared/src/model/web_socket.rs +++ b/shared/src/model/web_socket.rs @@ -1,6 +1,6 @@ use std::io; use bytes::Bytes; -use crate::model::{ConfigType, PlaylistUpdateState, StatusCheck}; +use crate::model::{ActiveUserConnectionChange, ConfigType, PlaylistUpdateState, StatusCheck}; use serde::{Deserialize, Serialize}; pub const PROTOCOL_VERSION: u8 = 1; @@ -79,7 +79,7 @@ pub enum ProtocolMessage { ServerError(String), StatusRequest(String), StatusResponse(StatusCheck), - ActiveUserResponse(usize, usize), // user_count, connection count + ActiveUserResponse(ActiveUserConnectionChange), ActiveProviderResponse(String, usize), ConfigChangeResponse(ConfigType), PlaylistUpdateResponse(PlaylistUpdateState),