diff --git a/README.md b/README.md index 766623281..b011b76d0 100644 --- a/README.md +++ b/README.md @@ -122,7 +122,7 @@ Messaging is Opt-In, you need to set the `notify_on` message types which are `telegram`, `rest` and `pushover.net` configurations are optional. `telegram` supports markdown generation for structured json messages. -`telegram` supports `message_thread_id` for groups chats. Simply put thread_id behind chat_id seperated by `:`. `':'` +`telegram` supports `message_thread_id` for group chats. Simply put thread_id behind chat_id separated by `:`. `':'` ```yaml messaging: diff --git a/backend/src/messaging.rs b/backend/src/messaging.rs index 204e6aa57..5411cac44 100644 --- a/backend/src/messaging.rs +++ b/backend/src/messaging.rs @@ -83,22 +83,20 @@ fn send_pushover_message(client: &Arc, msg: &str, messaging: &M } } -pub fn send_message_json(client: &Arc, kind: &MsgKind, cfg: Option<&MessagingConfig>, msg: &str) { +fn dispatch_send_message(client: &Arc, kind: MsgKind, cfg: Option<&MessagingConfig>, msg: &str, json: bool) { if let Some(messaging) = cfg { - if is_enabled(*kind, messaging) { - send_telegram_message(client, msg, messaging, true); + if is_enabled(kind, messaging) { + send_telegram_message(client, msg, messaging, json); send_http_post_request(client, msg, messaging); send_pushover_message(client, msg, messaging); } } } -pub fn send_message(client: &Arc, kind: &MsgKind, cfg: Option<&MessagingConfig>, msg: &str) { - if let Some(messaging) = cfg { - if is_enabled(*kind, messaging) { - send_telegram_message(client, msg, messaging, false); - send_http_post_request(client, msg, messaging); - send_pushover_message(client, msg, messaging); - } - } +pub fn send_message_json(client: &Arc, kind: MsgKind, cfg: Option<&MessagingConfig>, msg: &str) { + dispatch_send_message(client, kind, cfg, msg, true); +} + +pub fn send_message(client: &Arc, kind: MsgKind, cfg: Option<&MessagingConfig>, msg: &str) { + dispatch_send_message(client, kind, cfg, msg, false); } diff --git a/backend/src/model/stats.rs b/backend/src/model/stats.rs index 7a61673bb..bc21b0321 100644 --- a/backend/src/model/stats.rs +++ b/backend/src/model/stats.rs @@ -82,7 +82,7 @@ pub struct SourceStats { } impl SourceStats { - pub fn new(inputs: Vec, targets: Vec)-> Option { + pub fn try_new(inputs: Vec, targets: Vec) -> Option { if inputs.is_empty() || targets.is_empty() { None } else { diff --git a/backend/src/processing/playlist_watch.rs b/backend/src/processing/playlist_watch.rs index c87091657..650a207d5 100644 --- a/backend/src/processing/playlist_watch.rs +++ b/backend/src/processing/playlist_watch.rs @@ -73,7 +73,7 @@ fn handle_watch_notification(client: &Arc, cfg: &Config, added: let msg = serde_json::to_string_pretty(&changes).unwrap_or_else(|_| "Error: Failed to serialize watch changes".to_string()); info!("{}", &msg); - send_message(client, &MsgKind::Watch, cfg.messaging.as_ref(), &msg); + send_message(client, MsgKind::Watch, cfg.messaging.as_ref(), &msg); } } diff --git a/backend/src/processing/processor/playlist.rs b/backend/src/processing/processor/playlist.rs index 60de26ce9..41a61e05a 100644 --- a/backend/src/processing/processor/playlist.rs +++ b/backend/src/processing/processor/playlist.rs @@ -9,11 +9,12 @@ use std::sync::Arc; use std::thread; use tokio::sync::Mutex; -use crate::messaging::{send_message_json}; +use crate::api::model::{EventManager, EventMessage, PlaylistStorageState}; +use crate::messaging::send_message_json; use crate::model::Epg; -use crate::model::{ConfigTarget, ProcessTargets}; -use crate::model::{Mapping}; use crate::model::FetchedPlaylist; +use crate::model::Mapping; +use crate::model::{ConfigTarget, ProcessTargets}; use crate::model::{InputStats, PlaylistStats, SourceStats, TargetStats}; use crate::processing::parser::xmltv::flatten_tvguide; use crate::processing::playlist_watch::process_group_watch; @@ -34,7 +35,6 @@ use shared::model::{CounterModifier, FieldGetAccessor, FieldSetAccessor, InputTy PlaylistGroup, PlaylistItem, PlaylistUpdateState, ProcessingOrder, UUIDType, XtreamCluster}; use shared::utils::default_as_default; use std::time::Instant; -use crate::api::model::{EventManager, EventMessage, PlaylistStorageState}; fn is_valid(pli: &PlaylistItem, filter: &Filter) -> bool { let provider = ValueProvider { pli }; @@ -248,8 +248,8 @@ async fn playlist_download_from_input(client: &Arc, config: &Ar async fn process_source(client: Arc, cfg: Arc, source_idx: usize, user_targets: Arc, event_manager: Option>, - playlist_state: Option<&Arc> -)-> (Vec, Vec, Vec) { + playlist_state: Option<&Arc>, +) -> (Vec, Vec, Vec) { let sources = cfg.sources.load(); let mut errors = vec![]; let mut input_stats = HashMap::::new(); @@ -292,7 +292,7 @@ async fn process_source(client: Arc, cfg: Arc, sourc } let elapsed = start_time.elapsed().as_secs(); input_stats.insert(input_name.clone(), create_input_stat(group_count, channel_count, error_list.len(), - input.input_type, input_name, elapsed)); + input.input_type, input_name, elapsed)); } } if source_downloaded { @@ -340,7 +340,7 @@ fn create_input_stat(group_count: usize, channel_count: usize, error_count: usiz } async fn process_sources(client: Arc, config: &Arc, user_targets: Arc, - event_manager: Option>, playlist_state: Option<&Arc> + event_manager: Option>, playlist_state: Option<&Arc>, ) -> (Vec, Vec) { let mut handle_list = vec![]; let thread_num = config.config.load().threads; @@ -376,11 +376,11 @@ async fn process_sources(client: Arc, config: &Arc, let (input_stats, target_stats, mut res_errors) = process_source(Arc::clone(&http_client), cfg, index, usr_trgts, event_manager, playlist_state.as_ref()).await; shared_errors.lock().await.append(&mut res_errors); - if let Some(process_stats) = SourceStats::new(input_stats, target_stats) { + if let Some(process_stats) = SourceStats::try_new(input_stats, target_stats) { shared_stats.lock().await.push(process_stats); } }); - }, + } Err(err) => error!("Could not create runtime !!! {err}"), } }; @@ -392,7 +392,7 @@ async fn process_sources(client: Arc, config: &Arc, let (input_stats, target_stats, mut res_errors) = process_source(Arc::clone(&client), cfg, index, usr_trgts, event_manager, playlist_state).await; shared_errors.lock().await.append(&mut res_errors); - if let Some(process_stats) = SourceStats::new(input_stats, target_stats) { + if let Some(process_stats) = SourceStats::try_new(input_stats, target_stats) { shared_stats.lock().await.push(process_stats); } } @@ -478,7 +478,7 @@ async fn process_playlist_for_target(app_config: &AppConfig, stats: &mut HashMap, errors: &mut Vec, event_manager: Option>, - playlist_state: Option<&Arc> + playlist_state: Option<&Arc>, ) -> Result<(), Vec> { let pipe = get_processing_pipe(target); debug_if_enabled!("Processing order is {}", &target.processing_order); @@ -528,7 +528,7 @@ async fn process_playlist_for_target(app_config: &AppConfig, let mut flat_new_playlist = flatten_groups(new_playlist); step.tick("playlist merge"); - if sort_playlist(target, &mut flat_new_playlist) { + if sort_playlist(target, &mut flat_new_playlist) { step.tick("playlist sort"); } assign_channel_no_playlist(&mut flat_new_playlist); @@ -579,7 +579,7 @@ async fn process_epg(processed_fetched_playlists: &mut Vec>) } fn process_watch(cfg: &Config, client: &Arc, target: &ConfigTarget, new_playlist: &Vec) -> bool { - if let Some(watches) = &target.watch { + if let Some(watches) = &target.watch { if default_as_default().eq_ignore_ascii_case(&target.name) { error!("cant watch a target with no unique name"); } else { @@ -605,19 +605,29 @@ pub async fn exec_processing(client: Arc, app_config: Arc { + match serde_json::to_string(&serde_json::Value::Object( + serde_json::map::Map::from_iter([("stats".to_string(), val)]))) { + Ok(stats_msg) => { + // print stats + info!("{stats_msg}"); + // send stats + send_message_json(&client, MsgKind::Stats, messaging, stats_msg.as_str()); + } + Err(err) => error!("Failed to serialize playlist stats {err}"), + } + } + Err(err) => error!("Failed to serialize playlist stats {err}") } + // send errors if let Some(message) = get_errors_notify_message!(errors, 255) { if let Some(events) = event_manager { events.send_event(EventMessage::PlaylistUpdate(PlaylistUpdateState::Failure)); } if let Ok(error_msg) = serde_json::to_string(&serde_json::Value::Object(serde_json::map::Map::from_iter([("errors".to_string(), serde_json::Value::String(message))]))) { - send_message_json(&client, &MsgKind::Error, messaging, error_msg.as_str()); + send_message_json(&client, MsgKind::Error, messaging, error_msg.as_str()); } } else if let Some(events) = event_manager { events.send_event(EventMessage::PlaylistUpdate(PlaylistUpdateState::Success)); @@ -628,17 +638,17 @@ pub async fn exec_processing(client: Arc, app_config: Arc {}", first, second, strsim::jaro_winkler(first, second))); - // // println!("jaro {}", strsim::jaro(data.0, data.1)); - // // println!("levenhstein {}", strsim::levenshtein(data.0, data.1)); - // // println!("damerau_levenshtein {:?}", strsim::damerau_levenshtein(data.0, data.1)); - // // println!("osa distance {:?}", strsim::osa_distance(data.0, data.1)); - // // println!("sorensen dice {:?}", strsim::sorensen_dice(data.0, data.1)); - // } +// #[test] +// fn test_jaro_winkeler() { +// let data = [("yessport5", "heyessport5gold"), ("yessport5", "heyesport5gold")]; +// +// data.iter().for_each(|(first, second)| +// println!("jaro_winkler {} = {} => {}", first, second, strsim::jaro_winkler(first, second))); +// // println!("jaro {}", strsim::jaro(data.0, data.1)); +// // println!("levenhstein {}", strsim::levenshtein(data.0, data.1)); +// // println!("damerau_levenshtein {:?}", strsim::damerau_levenshtein(data.0, data.1)); +// // println!("osa distance {:?}", strsim::osa_distance(data.0, data.1)); +// // println!("sorensen dice {:?}", strsim::sorensen_dice(data.0, data.1)); +// } // } \ No newline at end of file diff --git a/backend/src/utils/network/xtream.rs b/backend/src/utils/network/xtream.rs index 5f07bf256..5486a7dec 100644 --- a/backend/src/utils/network/xtream.rs +++ b/backend/src/utils/network/xtream.rs @@ -157,7 +157,7 @@ async fn xtream_login(cfg: &Config, client: &Arc, input: &Input if let Ok(cur_status) = ProxyUserStatus::from_str(&status) { if !matches!(cur_status, ProxyUserStatus::Active | ProxyUserStatus::Trial) { warn!("User status for user {username} is {cur_status:?}"); - send_message(client, &MsgKind::Info, cfg.messaging.as_ref(), &format!("User status for user {username} is {cur_status:?}")); + send_message(client, MsgKind::Info, cfg.messaging.as_ref(), &format!("User status for user {username} is {cur_status:?}")); } } } @@ -175,11 +175,11 @@ async fn xtream_login(cfg: &Config, client: &Arc, input: &Input let datetime = DateTime::from_timestamp(expiration_timestamp, 0).unwrap(); let formatted = datetime.format("%Y-%m-%d %H:%M:%S").to_string(); warn!("User account for user {username} expires {formatted}"); - send_message(client, &MsgKind::Info, cfg.messaging.as_ref(), &format!("User account for user {username} expires {formatted}")); + send_message(client, MsgKind::Info, cfg.messaging.as_ref(), &format!("User account for user {username} expires {formatted}")); } } else { warn!("User account for user {username} is expired"); - send_message(client, &MsgKind::Info, cfg.messaging.as_ref(), &format!("User account for user {username} is expired")); + send_message(client, MsgKind::Info, cfg.messaging.as_ref(), &format!("User account for user {username} is expired")); } } } diff --git a/shared/src/utils/json_utils.rs b/shared/src/utils/json_utils.rs index 6eeac2818..06ffff0c8 100644 --- a/shared/src/utils/json_utils.rs +++ b/shared/src/utils/json_utils.rs @@ -186,20 +186,27 @@ fn json_to_markdown(value: &Value) -> String { fn format_value(v: &Value, indent: usize) -> String { let pad = " ".repeat(indent); match v { - Value::Object(map) => map.iter() - .map(|(k, v)| { - let formatted = format_value(v, indent + 1); - let key = humanize_snake_case(k); - if v.is_object() || v.is_array() { - format!("{pad}*{key}:*\n{formatted}") - } else { - format!("{pad}*{key}:* {formatted}") - } - }) - .collect::>() - .join("\n"), + Value::Object(map) => { + let mut entries: Vec<_> = map.iter().collect(); + entries.sort_by_key(|(k, _)| *k); + entries.into_iter() + .map(|(k, v)| { + let formatted = format_value(v, indent + 1); + let key = escape_markdown_v2(&humanize_snake_case(k)); + if v.is_object() || v.is_array() { + format!("{pad}*{key}:*\n{formatted}") + } else { + format!("{pad}*{key}:* {formatted}") + } + }) + .collect::>() + .join("\n") + }, Value::Array(arr) => arr.iter() - .map(|v| format!("{pad}\\- {}", format_value(v, indent + 1).trim())) + .map(|v| { + let dash_pad = if indent > 0 { " ".repeat(indent - 1) } else { "".to_string() }; + format!("{dash_pad}\\- {}", format_value(v, indent + 1).trim()) + }) .collect::>() .join("\n"), Value::String(s) => escape_markdown_v2(s),