From 7c7c467f64e47dcd1defa9ebff2ef6c301b93c02 Mon Sep 17 00:00:00 2001 From: euzu Date: Sat, 25 Oct 2025 20:40:58 +0200 Subject: [PATCH] Stream table new columns --- Cargo.lock | 1 + backend/src/api/api_utils.rs | 45 +++++---- backend/src/api/model/active_user_manager.rs | 4 +- .../api/model/streams/active_client_stream.rs | 19 ++-- backend/src/repository/bplustree.rs | 98 +++++++++++++++++++ backend/src/repository/xtream_repository.rs | 2 + backend/src/utils/geoip.rs | 84 ++++++++++++++++ backend/src/utils/mod.rs | 1 + frontend/public/assets/i18n/en.json | 6 +- .../components/dashboard/_streams_view.scss | 4 + .../app/components/dashboard/streams_table.rs | 66 +++++++++++-- shared/Cargo.toml | 1 + shared/src/model/config/web_ui.rs | 1 + shared/src/model/stream_info.rs | 11 ++- shared/src/utils/string_utils.rs | 10 ++ shared/src/utils/time_utils.rs | 15 +-- 16 files changed, 319 insertions(+), 49 deletions(-) create mode 100644 backend/src/utils/geoip.rs diff --git a/Cargo.lock b/Cargo.lock index 054477ff3..8a7890b16 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3868,6 +3868,7 @@ dependencies = [ "enum-iterator", "fastrand", "indexmap", + "js-sys", "log", "path-clean", "pest", diff --git a/backend/src/api/api_utils.rs b/backend/src/api/api_utils.rs index c9633a2e2..4ef2cf2da 100644 --- a/backend/src/api/api_utils.rs +++ b/backend/src/api/api_utils.rs @@ -329,7 +329,7 @@ enum ProviderStreamState { pub struct StreamDetails { pub stream: Option, stream_info: ProviderStreamInfo, - pub input_name: Option, + pub provider_name: Option, pub grace_period_millis: u64, pub reconnect_flag: Option>, pub provider_connection_guard: Option>, @@ -340,7 +340,7 @@ impl StreamDetails { Self { stream: Some(stream), stream_info: None, - input_name: None, + provider_name: None, grace_period_millis: default_grace_period_millis(), reconnect_flag: None, provider_connection_guard: None, @@ -484,6 +484,11 @@ async fn create_stream_response_details( &streaming_strategy.provider_stream_state, config_grace_period_millis, ); + let provider_name = streaming_strategy + .provider_connection_guard + .as_ref() + .and_then(|guard| guard.get_provider_name()); + match streaming_strategy.provider_stream_state { // custom stream means we display our own stream like connection exhausted, channel-unavailable... ProviderStreamState::Custom(provider_stream) => { @@ -491,7 +496,7 @@ async fn create_stream_response_details( StreamDetails { stream, stream_info, - input_name: None, + provider_name: provider_name.clone(), grace_period_millis, reconnect_flag: None, provider_connection_guard: streaming_strategy.provider_connection_guard.clone(), @@ -525,14 +530,6 @@ async fn create_stream_response_details( ((None, None), None) }; - // if we have no stream, we should release the provider - if stream.is_none() { - if let Some(guard) = streaming_strategy.provider_connection_guard.take() { - drop(guard); - } - error!("Cant open stream {}", sanitize_sensitive_info(&request_url)); - } - if log_enabled!(log::Level::Debug) { if let Some((headers, status_code, response_url)) = stream_info.as_ref() { debug!( @@ -546,10 +543,18 @@ async fn create_stream_response_details( } } + // if we have no stream, we should release the provider + if stream.is_none() { + if let Some(guard) = streaming_strategy.provider_connection_guard.take() { + drop(guard); + } + error!("Cant open stream {}", sanitize_sensitive_info(&request_url)); + } + StreamDetails { stream, stream_info, - input_name: provider_name, + provider_name, grace_period_millis, reconnect_flag, provider_connection_guard: streaming_strategy.provider_connection_guard.take(), @@ -770,7 +775,7 @@ pub async fn force_provider_stream_response( .await; stream_channel.shared = share_stream; let stream = - ActiveClientStream::new(stream_details, app_state, user, connection_permission, addr, stream_channel) + ActiveClientStream::new(stream_details, app_state, user, connection_permission, addr, stream_channel, req_headers) .await; let (status_code, header_map) = @@ -834,7 +839,7 @@ pub async fn stream_response( 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, stream_channel.clone()).await + shared_stream_response(app_state, stream_url, addr, user, connection_permission, stream_channel.clone(), req_headers).await { return value.into_response(); } @@ -860,10 +865,7 @@ pub async fn stream_response( .stream_info .as_ref() .map(|(h, sc, response_url)| (h.clone(), *sc, response_url.clone())); - let provider_name = stream_details - .provider_connection_guard - .as_ref() - .and_then(|guard| guard.get_provider_name()); + let provider_name = stream_details.provider_name.clone(); let provider_guard = if share_stream { stream_details.provider_connection_guard.take() @@ -872,7 +874,7 @@ pub async fn stream_response( }; stream_channel.shared = share_stream; let stream = - ActiveClientStream::new(stream_details, app_state, user, connection_permission, addr, stream_channel) + ActiveClientStream::new(stream_details, app_state, user, connection_permission, addr, stream_channel, req_headers) .await; let stream_resp = if share_stream { debug_if_enabled!( @@ -986,7 +988,8 @@ async fn shared_stream_response( addr: &str, user: &ProxyUserCredentials, connect_permission: UserConnectionPermission, - mut stream_channel: StreamChannel + mut stream_channel: StreamChannel, + req_headers: &HeaderMap, ) -> Option { if let Some(stream) = SharedStreamManager::subscribe_shared_stream(app_state, stream_url, Some(addr)).await @@ -1007,7 +1010,7 @@ async fn shared_stream_response( let stream_details = StreamDetails::from_stream(stream); stream_channel.shared = true; let stream = - ActiveClientStream::new(stream_details, app_state, user, connect_permission, addr, stream_channel) + ActiveClientStream::new(stream_details, app_state, user, connect_permission, addr, stream_channel, req_headers) .await .boxed(); let mut response = axum::response::Response::builder().status(status_code); diff --git a/backend/src/api/model/active_user_manager.rs b/backend/src/api/model/active_user_manager.rs index 45d26d478..27177e5ff 100644 --- a/backend/src/api/model/active_user_manager.rs +++ b/backend/src/api/model/active_user_manager.rs @@ -1,3 +1,4 @@ +use std::borrow::Cow; use crate::api::model::ActiveProviderManager; use crate::api::model::SharedStreamManager; use crate::model::Config; @@ -280,12 +281,13 @@ impl ActiveUserManager { Self::get_active_connections(&self.user).await } - pub async fn add_connection(&self, username: &str, max_connections: u32, addr: &str, provider: &str, stream_channel: StreamChannel) -> UserConnectionGuard { + pub async fn add_connection(&self, username: &str, max_connections: u32, addr: &str, provider: &str, stream_channel: StreamChannel, user_agent: Cow<'_, str>) -> UserConnectionGuard { let stream_info = StreamInfo::new( username, addr, provider, stream_channel, + user_agent.to_string(), ); { let mut user_map = self.user.write().await; diff --git a/backend/src/api/model/streams/active_client_stream.rs b/backend/src/api/model/streams/active_client_stream.rs index cb2613fd6..437c83952 100644 --- a/backend/src/api/model/streams/active_client_stream.rs +++ b/backend/src/api/model/streams/active_client_stream.rs @@ -15,6 +15,8 @@ use std::pin::Pin; use std::sync::atomic::AtomicU8; use std::sync::{Arc}; use std::task::{Poll}; +use axum::http::header::USER_AGENT; +use axum::http::HeaderMap; use futures::task::AtomicWaker; const INNER_STREAM: u8 = 0_u8; @@ -39,19 +41,16 @@ impl ActiveClientStream { user: &ProxyUserCredentials, connection_permission: UserConnectionPermission, addr: &str, - stream_channel: StreamChannel) -> Self { + stream_channel: StreamChannel, + req_headers: &HeaderMap) -> 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 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 provider_name = stream_details.provider_name.as_ref().map_or_else(String::new, ToString::to_string); + let user_agent = req_headers.get(USER_AGENT).map(|h| String::from_utf8_lossy(h.as_bytes())).unwrap_or_default(); + let user_connection_guard = Some(app_state.active_users.add_connection(username, user.max_connections, addr, &provider_name, stream_channel, user_agent).await); let cfg = &app_state.app_config; let waker = Some(Arc::new(AtomicWaker::new())); let waker_clone = waker.clone(); @@ -106,8 +105,8 @@ impl ActiveClientStream { let active_provider = Arc::clone(&app_state.active_provider); let shared_stream_manager = Arc::clone(&app_state.shared_stream_manager); - let provider_grace_check = if stream_details.has_grace_period() && stream_details.input_name.is_some() { - let provider_name = stream_details.input_name.as_deref().unwrap_or_default().to_string(); + let provider_grace_check = if stream_details.has_grace_period() && stream_details.provider_name.is_some() { + let provider_name = stream_details.provider_name.as_ref().map_or_else(String::new, ToString::to_string); Some(provider_name) } else { None diff --git a/backend/src/repository/bplustree.rs b/backend/src/repository/bplustree.rs index 06a845e1e..c9bdeaff7 100644 --- a/backend/src/repository/bplustree.rs +++ b/backend/src/repository/bplustree.rs @@ -58,6 +58,50 @@ where } +fn query_tree_le(file: &mut R, key: &K) -> Option +where + K: Ord + Serialize + for<'de> Deserialize<'de> + Clone, + V: Serialize + for<'de> Deserialize<'de> + Clone, +{ + let mut offset = 0; + let mut buffer = vec![0u8; BLOCK_SIZE]; + loop { + match BPlusTreeNode::::deserialize_from_block(file, &mut buffer, offset, false) { + Ok((node, pointers)) => { + if node.is_leaf { + let idx = get_entry_index_upper_bound::(&node.keys, key); + if idx == 0 { + return None; + } else { + return node.values.get(idx - 1).cloned(); + } + } + let child_idx = get_entry_index_upper_bound::(&node.keys, key); + if let Some(child_offsets) = pointers { + if let Some(child_offset) = child_offsets.get(child_idx) { + offset = *child_offset; + } else { + // defensive: if out of bounds try last pointer + if let Some(last) = child_offsets.last() { + offset = *last; + } else { + return None; + } + } + } else { + return None; + } + } + Err(err) => { + error!("Failed to read id tree from file {err}"); + return None; + } + } + } +} + + + #[derive(Serialize, Deserialize, Debug, Clone)] struct BPlusTreeNode { keys: Vec, @@ -193,6 +237,32 @@ where } } + /// Find the largest key <= `key` in this subtree. + /// Returns a reference to (key, value) if found (only valid for leaf entries). + fn find_le(&self, key: &K) -> Option<(&K, &V)> { + if self.is_leaf { + // find index of first key > key, then step one back + let idx = self.get_entry_index_upper_bound(key); + if idx == 0 { + None + } else { + let i = idx - 1; + // safe: leaf guarantees values.len() == keys.len() + Some((&self.keys[i], &self.values[i])) + } + } else { + // descend into the appropriate child (child index = upper_bound) + let child_idx = self.get_entry_index_upper_bound(key); + // child_idx can be equal to children.len() if key > all keys; children.get handles that + if let Some(child) = self.children.get(child_idx) { + child.find_le(key) + } else { + // fallback: if child_idx is out of bounds, try last child (defensive) + self.children.last().and_then(|c| c.find_le(key)) + } + } + } + pub fn traverse(&self, visit: &mut F) where F: FnMut(&Vec, &Vec), @@ -506,6 +576,15 @@ where Ok(Self::new_with_root(root)) } + /// Find the largest key <= `key` in the in-memory tree and return references to (key, value). + pub fn find_le(&self, key: &K) -> Option<(&K, &V)> { + // empty tree + if self.root.keys.is_empty() && self.root.is_leaf && self.root.values.is_empty() { + return None; + } + self.root.find_le(key) + } + pub fn traverse(&self, mut visit: F) where F: FnMut(&Vec, &Vec), @@ -604,6 +683,19 @@ where query_tree(&mut self.file, key) } + /// On-disk: find largest key <= `key` and return owned V (cloned/deserialized) + pub fn query_le(&mut self, key: &K) -> Option { + // use the same buffer/reader pattern as query() + // we need a mutable reader over the inner BufReader + let file = &mut self.file; + // Seek to start to be safe + if file.seek(SeekFrom::Start(0)).is_err() { + // if seek fails, still try to query — but bail out with None + return None; + } + query_tree_le(file, key) + } + // pub fn traverse(&mut self, mut visit: F) // where // F: FnMut(&Vec, &Vec), @@ -681,6 +773,12 @@ where } } } + + /// On-disk update helper: find largest key <= `key`. + pub fn query_le(&mut self, key: &K) -> Option { + let mut reader = utils::file_reader(&mut self.file); + query_tree_le(&mut reader, key) + } } pub struct BPlusTreeIterator<'a, K, V> { diff --git a/backend/src/repository/xtream_repository.rs b/backend/src/repository/xtream_repository.rs index c451445e2..878ae1e7f 100644 --- a/backend/src/repository/xtream_repository.rs +++ b/backend/src/repository/xtream_repository.rs @@ -440,6 +440,8 @@ pub async fn xtream_get_item_for_stream_id( PlaylistItemType::Series => { if let Ok(mut item) = xtream_read_series_item_for_stream_id(app_config, mapping.parent_virtual_id, &storage_path) { item.provider_id = mapping.provider_id; + + Ok(item) } else { xtream_read_item_for_stream_id(app_config, virtual_id, &storage_path, XtreamCluster::Series) diff --git a/backend/src/utils/geoip.rs b/backend/src/utils/geoip.rs new file mode 100644 index 000000000..f3ea4f25e --- /dev/null +++ b/backend/src/utils/geoip.rs @@ -0,0 +1,84 @@ +use std::io; +use std::io::{BufRead}; +use std::net::Ipv4Addr; +use std::path::Path; +use serde::{Serialize, Deserialize}; +use crate::repository::bplustree::BPlusTree; + +fn ipv4_to_u32(ip: &str) -> Option { + ip.parse::().ok().map(|a| u32::from(a)) +} + +#[derive(Serialize, Deserialize)] +pub struct GeoIp { + tree: BPlusTree, +} + +impl GeoIp { + + pub fn load(path: &Path) -> io::Result { + let tree = BPlusTree::load(path)?; + Ok(Self { tree }) + } + + pub fn new() -> Self { + Self { tree: BPlusTree::new() } + } + + pub fn import_ipv4_from_csv(&mut self, mut reader: impl BufRead, db_path: &Path) -> std::io::Result { + let mut buf = String::new(); + + while reader.read_line(&mut buf)? > 0 { + let line = buf.trim(); + if line.is_empty() || line.starts_with('#') { continue; } + + let parts: Vec<&str> = line.split(',').collect(); + if parts.len() != 3 { continue; } + + if let (Some(start), Some(end)) = (ipv4_to_u32(parts[0]), ipv4_to_u32(parts[1])) { + let cc = parts[2].trim().to_string(); + self.tree.insert(start, (end, cc)); + } + buf.clear(); + } + self.tree.store(db_path) + } + + pub fn lookup(&self, ip_str: &str) -> Option { + let ip = ipv4_to_u32(ip_str)?; + if let Some((_, (end, cc))) = self.tree.find_le(&ip) { + if ip <= *end { + return Some(cc.to_string()); + } + } + None + } +} + +#[cfg(test)] +mod test { + //https://github.com/datasets/geoip2-ipv4/blob/main/data/geoip2-ipv4.csv + + use std::fs::File; + use std::io::BufReader; + use std::path::PathBuf; + use crate::utils::geoip::GeoIp; + + #[test] + pub fn test_csv() { + let db_file = PathBuf::from("/projects/m3u-test/asn-country-ipv4.db"); + let source = PathBuf::from("/projects/m3u-test/asn-country-ipv4.csv"); + let file = File::open(source).expect("Could not open csv file"); + let reader = BufReader::new(file); + let mut geo_ip = GeoIp::new(); + let _ = geo_ip.import_ipv4_from_csv(reader, &db_file).expect("Could not import csv"); + + let geo_ip = GeoIp::load(&db_file).expect("Failed to load geoip db"); + if let Some(cc) = geo_ip.lookup("72.13.24.23") { + assert_eq!(cc, "US"); + } else { + assert!(false); + } + + } +} \ No newline at end of file diff --git a/backend/src/utils/mod.rs b/backend/src/utils/mod.rs index a99877bb3..d0477a11d 100644 --- a/backend/src/utils/mod.rs +++ b/backend/src/utils/mod.rs @@ -9,6 +9,7 @@ mod trakt; mod json_utils; mod bincode_utils; mod telegram; +mod geoip; pub use self::bincode_utils::*; pub use self::logging::*; diff --git a/frontend/public/assets/i18n/en.json b/frontend/public/assets/i18n/en.json index 6c639ab36..339dbb62f 100644 --- a/frontend/public/assets/i18n/en.json +++ b/frontend/public/assets/i18n/en.json @@ -290,7 +290,11 @@ "GROUP": "Group", "CLIENT_IP": "Client IP", "STREAM_ID": "Stream Id", - "SHARED": "Shared" + "SHARED": "Shared", + "CLUSTER": "Type", + "USER_AGENT": "Player", + "FORMAT": "Format", + "DURATION": "Duration" }, "TITLE": { "USER_BOUQUET_EDITOR": "User group editor" diff --git a/frontend/scss/app/components/dashboard/_streams_view.scss b/frontend/scss/app/components/dashboard/_streams_view.scss index 8ad4295eb..d8439096c 100644 --- a/frontend/scss/app/components/dashboard/_streams_view.scss +++ b/frontend/scss/app/components/dashboard/_streams_view.scss @@ -18,4 +18,8 @@ gap: var(--gap-larger); overflow: auto; } +} + +.tp__stream-table__duration { + letter-spacing: 2px; } \ No newline at end of file diff --git a/frontend/src/app/components/dashboard/streams_table.rs b/frontend/src/app/components/dashboard/streams_table.rs index 53b87cd65..0f87136e9 100644 --- a/frontend/src/app/components/dashboard/streams_table.rs +++ b/frontend/src/app/components/dashboard/streams_table.rs @@ -1,3 +1,4 @@ +use std::borrow::Cow; use crate::app::components::menu_item::MenuItem; use crate::app::components::popup_menu::PopupMenu; use crate::app::components::{AppIcon, Table, TableDefinition, ToggleSwitch}; @@ -7,20 +8,60 @@ use shared::model::{SortOrder, StreamInfo}; use std::fmt::Display; use std::rc::Rc; use std::str::FromStr; +use gloo_timers::callback::Interval; +use gloo_utils::window; +use log::debug; +use wasm_bindgen::JsCast; +use web_sys::Element; use yew::prelude::*; use yew_i18n::use_translation; +use shared::utils::current_time_secs; -const HEADERS: [&str; 8] = [ +const HEADERS: [&str; 11] = [ "LABEL.EMPTY", "LABEL.USERNAME", "LABEL.STREAM_ID", + "LABEL.CLUSTER", "LABEL.CHANNEL", "LABEL.GROUP", "LABEL.CLIENT_IP", "LABEL.PROVIDER", - "LABEL.SHARED" + "LABEL.SHARED", + "LABEL.USER_AGENT", + "LABEL.DURATION" ]; +fn strip_port<'a>(input: &'a str) -> Cow<'a, str> { + if let Some(pos) = input.find(':') { + Cow::Owned(input[..pos].to_string()) + } else { + Cow::Borrowed(input) + } +} + +pub fn format_duration(seconds: u64) -> String { + let hours = seconds / 3600; + let minutes = (seconds % 3600) / 60; + let seconds = seconds % 60; + format!("{hours:02}:{minutes:02}:{seconds:02}") +} + +fn update_timestamps() { + let window = window(); + let document = window.document().unwrap(); + let spans = document.query_selector_all("span[data-ts]").unwrap(); + for i in 0..spans.length() { + if let Some(node) = spans.item(i) { + let el: Element = node.dyn_into().unwrap(); + if let Some(ts_str) = el.get_attribute("data-ts") { + if let Ok(ts) = ts_str.parse::() { + el.set_inner_html(&format_duration(current_time_secs() - ts)); + } + } + } + } +} + #[derive(Properties, PartialEq, Clone)] pub struct StreamsTableProps { pub streams: Option>>, @@ -34,6 +75,14 @@ pub fn StreamsTable(props: &StreamsTableProps) -> Html { let popup_is_open = use_state(|| false); let selected_dto = use_state(|| None::>); + + use_effect_with((), move |_| { + Interval::new(1000, || { + update_timestamps(); + }).forget(); + }); + + let handle_popup_close = { let set_is_open = popup_is_open.clone(); Callback::from(move |()| { @@ -91,11 +140,14 @@ pub fn StreamsTable(props: &StreamsTableProps) -> Html { { 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()}, - 7 => html! { }, + 3 => html! {dto.channel.cluster}, + 4 => html! {dto.channel.title.as_str()}, + 5 => html! {dto.channel.group.as_str()}, + 6 => html! { strip_port(&dto.addr)}, + 7 => html! {dto.provider.as_str()}, + 8 => html! { }, + 9 => html! { dto.user_agent.as_str() }, + 10 => html! { {format_duration(dto.ts)} }, _ => html! {""}, } }) diff --git a/shared/Cargo.toml b/shared/Cargo.toml index bc1e091ec..bb30e6594 100644 --- a/shared/Cargo.toml +++ b/shared/Cargo.toml @@ -22,3 +22,4 @@ zeroize = "1" chrono = "0.4.42" bytes = "1" ciborium = "0.2.2" +js-sys = "0.3.81" diff --git a/shared/src/model/config/web_ui.rs b/shared/src/model/config/web_ui.rs index c0aac7aeb..c16f21b08 100644 --- a/shared/src/model/config/web_ui.rs +++ b/shared/src/model/config/web_ui.rs @@ -3,6 +3,7 @@ use crate::model::WebAuthConfigDto; use crate::utils::{default_as_true, is_blank_optional_string}; const RESERVED_PATHS: &[&str] = &[ + "cvs", "live", "movie", "series", diff --git a/shared/src/model/stream_info.rs b/shared/src/model/stream_info.rs index 0edd1d460..2b0ed889d 100644 --- a/shared/src/model/stream_info.rs +++ b/shared/src/model/stream_info.rs @@ -1,5 +1,6 @@ use serde::{Deserialize, Serialize}; use crate::model::{M3uPlaylistItem, PlaylistEntry, PlaylistItemType, XtreamCluster, XtreamPlaylistItem}; +use crate::utils::{current_time_secs, StringExt}; #[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] pub struct StreamChannel { @@ -20,7 +21,7 @@ impl XtreamPlaylistItem { item_type: self.item_type, cluster: self.xtream_cluster, group: self.group.clone(), - title: self.title.clone(), + title: String::longest(self.title.as_str(), self.name.as_str()).to_string(), url: self.url.clone(), shared: false, } @@ -35,7 +36,7 @@ impl M3uPlaylistItem { item_type: self.item_type, cluster: XtreamCluster::try_from(self.item_type).unwrap_or(XtreamCluster::Live), group: self.group.clone(), - title: self.title.clone(), + title: String::longest(self.title.as_str(), self.name.as_str()).to_string(), url: self.url.clone(), shared: false, } @@ -48,15 +49,19 @@ pub struct StreamInfo { pub channel: StreamChannel, pub provider: String, pub addr: String, + pub user_agent: String, + pub ts: u64, } impl StreamInfo { - pub fn new(username: &str, addr: &str, provider: &str, stream_channel: StreamChannel) -> Self { + pub fn new(username: &str, addr: &str, provider: &str, stream_channel: StreamChannel, user_agent: String) -> Self { Self { username: username.to_string(), channel: stream_channel, provider: provider.to_string(), addr: addr.to_string(), + user_agent, + ts: current_time_secs(), } } } \ No newline at end of file diff --git a/shared/src/utils/string_utils.rs b/shared/src/utils/string_utils.rs index 4c6f333d4..67a957a80 100644 --- a/shared/src/utils/string_utils.rs +++ b/shared/src/utils/string_utils.rs @@ -127,6 +127,16 @@ pub fn humanize_snake_case(s: &str) -> String { result } +pub trait StringExt { + fn longest<'a>(a: &'a str, b: &'a str) -> &'a str; +} + +impl StringExt for String { + fn longest<'a>(a: &'a str, b: &'a str) -> &'a str { + if a.len() >= b.len() { a } else { b } + } +} + #[cfg(test)] mod test { use std::collections::HashSet; diff --git a/shared/src/utils/time_utils.rs b/shared/src/utils/time_utils.rs index 0fb480747..87267eb5a 100644 --- a/shared/src/utils/time_utils.rs +++ b/shared/src/utils/time_utils.rs @@ -1,9 +1,12 @@ -use std::time::{SystemTime, UNIX_EPOCH}; -use chrono::{DateTime}; - +#[cfg(target_arch = "wasm32")] pub fn current_time_secs() -> u64 { - SystemTime::now() - .duration_since(UNIX_EPOCH) + (js_sys::Date::now() / 1000.0) as u64 +} + +#[cfg(not(target_arch = "wasm32"))] +pub fn current_time_secs() -> u64 { + std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) .unwrap_or_default() .as_secs() } @@ -15,7 +18,7 @@ pub fn unix_ts_to_str(ts: i64) -> Option { } else { ts }; - DateTime::from_timestamp(normalized_ts, 0).map(|dt| dt.format("%d.%m.%Y").to_string()) + chrono::DateTime::from_timestamp(normalized_ts, 0).map(|dt| dt.format("%d.%m.%Y").to_string()) } else { None }