From 86d840939177e714045ff291ff88a437b748cc00 Mon Sep 17 00:00:00 2001 From: euzu Date: Tue, 21 Oct 2025 23:13:02 +0200 Subject: [PATCH 1/3] Telegram thread_id support --- backend/src/utils/mod.rs | 1 - backend/src/utils/telegram.rs | 2 +- 2 files changed, 1 insertion(+), 2 deletions(-) diff --git a/backend/src/utils/mod.rs b/backend/src/utils/mod.rs index 5036ce991..a99877bb3 100644 --- a/backend/src/utils/mod.rs +++ b/backend/src/utils/mod.rs @@ -15,7 +15,6 @@ pub use self::logging::*; pub use self::trakt::*; pub use self::telegram::*; - #[macro_export] macro_rules! debug_if_enabled { ($fmt:expr, $( $args:expr ),*) => { diff --git a/backend/src/utils/telegram.rs b/backend/src/utils/telegram.rs index 661a1bda5..d0d26a66c 100644 --- a/backend/src/utils/telegram.rs +++ b/backend/src/utils/telegram.rs @@ -39,7 +39,7 @@ fn get_send_message_parse_mode_str(mode: &SendMessageParseMode) -> &'static str } #[derive(Debug, serde::Serialize)] -pub struct RequestObj { +struct RequestObj { pub chat_id: String, #[serde(skip_serializing_if = "Option::is_none")] pub message_thread_id: Option, From 8245bc4c85e4c8a0425543cec80aaefaf1f31913 Mon Sep 17 00:00:00 2001 From: euzu Date: Wed, 22 Oct 2025 11:04:23 +0200 Subject: [PATCH 2/3] Telegram markdown support for structured json messages --- CHANGELOG.md | 1 + README.md | 4 ++ backend/src/api/endpoints/m3u_api.rs | 5 +- backend/src/messaging.rs | 42 ++++++++++++---- backend/src/model/config/messaging.rs | 3 ++ backend/src/model/stats.rs | 8 ++- backend/src/processing/playlist_watch.rs | 34 +++++++------ backend/src/processing/processor/playlist.rs | 16 +++--- backend/src/utils/telegram.rs | 8 +-- frontend/public/assets/i18n/en.json | 1 + .../config/messaging_config_view.rs | 14 +++--- shared/src/model/config/messaging.rs | 2 + shared/src/utils/json_utils.rs | 49 +++++++++++++++++++ shared/src/utils/string_utils.rs | 21 ++++++++ 14 files changed, 160 insertions(+), 48 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 8fb3bdc67..86d9d9616 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -11,6 +11,7 @@ Placing playlist into memory causes more RAM usage but reduces disk access. Output filters are applied after all transformations have been performed, therefore, all filter contents must refer to the final state of the playlist. - 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. # 3.1.7 (2025-10-10) - Added Dark/Bright theme switch diff --git a/README.md b/README.md index bbdbb8080..766623281 100644 --- a/README.md +++ b/README.md @@ -121,6 +121,9 @@ 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 `:`. `':'` + ```yaml messaging: notify_on: @@ -128,6 +131,7 @@ messaging: - stats - error telegram: + markdown: true bot_token: '' chat_ids: - '' diff --git a/backend/src/api/endpoints/m3u_api.rs b/backend/src/api/endpoints/m3u_api.rs index d950206a2..27621c81b 100644 --- a/backend/src/api/endpoints/m3u_api.rs +++ b/backend/src/api/endpoints/m3u_api.rs @@ -16,10 +16,7 @@ use axum::response::IntoResponse; use bytes::Bytes; use futures::stream; use log::{debug, error}; -use shared::model::{ - FieldGetAccessor, PlaylistEntry, PlaylistItemType, TargetType, UserConnectionPermission, - XtreamCluster, -}; +use shared::model::{FieldGetAccessor, PlaylistEntry, PlaylistItemType, TargetType, UserConnectionPermission, XtreamCluster}; use shared::utils::{concat_path, extract_extension_from_url, sanitize_sensitive_info, HLS_EXT}; use std::sync::Arc; diff --git a/backend/src/messaging.rs b/backend/src/messaging.rs index bb6dac6b8..204e6aa57 100644 --- a/backend/src/messaging.rs +++ b/backend/src/messaging.rs @@ -1,9 +1,11 @@ -use std::sync::Arc; -use crate::model::{MessagingConfig}; +use std::borrow::Cow; +use crate::model::MessagingConfig; +use crate::utils::{telegram_create_instance, telegram_send_message, SendMessageOption, SendMessageParseMode}; use log::{debug, error}; -use reqwest::{header}; +use reqwest::header; use shared::model::MsgKind; -use crate::utils::{telegram_create_instance, telegram_send_message}; +use shared::utils::json_str_to_markdown; +use std::sync::Arc; fn is_enabled(kind: MsgKind, cfg: &MessagingConfig) -> bool { cfg.notify_on.contains(&kind) @@ -29,12 +31,24 @@ fn send_http_post_request(client: &Arc, msg: &str, messaging: & } } -fn send_telegram_message(client: &Arc, msg: &str, messaging: &MessagingConfig) { +fn send_telegram_message(client: &Arc, msg: &str, messaging: &MessagingConfig, json: bool) { // TODO use proxy settings if let Some(telegram) = &messaging.telegram { + let (message, options) = { + if json && telegram.markdown { + if let Ok(md) = json_str_to_markdown(msg) { + (Cow::Owned(md), Some(SendMessageOption { parse_mode: SendMessageParseMode::MarkdownV2 })) + } else { + (Cow::Borrowed(msg), None) + } + } else { + (Cow::Borrowed(msg), None) + } + }; + for chat_id in &telegram.chat_ids { let bot = telegram_create_instance(&telegram.bot_token, chat_id); - telegram_send_message(client, &bot, msg, None); + telegram_send_message(client, &bot, &message, options.as_ref()); } } } @@ -62,17 +76,27 @@ fn send_pushover_message(client: &Arc, msg: &str, messaging: &M } else { error!("Failed to send text message to PUSHOVER, status code {}", response.status()); } - }, + } Err(e) => error!("Text message wasn't sent to PUSHOVER api because of: {e}"), } }); } } -pub fn send_message(client: &Arc, kind: &MsgKind, cfg: Option<&MessagingConfig>, msg: &str) { +pub fn send_message_json(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); + send_telegram_message(client, msg, messaging, true); + 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); } diff --git a/backend/src/model/config/messaging.rs b/backend/src/model/config/messaging.rs index f6812495e..a9855b55b 100644 --- a/backend/src/model/config/messaging.rs +++ b/backend/src/model/config/messaging.rs @@ -5,6 +5,7 @@ use crate::model::macros; pub struct TelegramMessagingConfig { pub bot_token: String, pub chat_ids: Vec, + pub markdown: bool, } macros::from_impl!(TelegramMessagingConfig); @@ -13,6 +14,7 @@ impl From<&TelegramMessagingConfigDto> for TelegramMessagingConfig { Self { bot_token: dto.bot_token.clone(), chat_ids: dto.chat_ids.clone(), + markdown: dto.markdown, } } } @@ -22,6 +24,7 @@ impl From<&TelegramMessagingConfig> for TelegramMessagingConfigDto { Self { bot_token: instance.bot_token.clone(), chat_ids: instance.chat_ids.clone(), + markdown: instance.markdown, } } } diff --git a/backend/src/model/stats.rs b/backend/src/model/stats.rs index 9488b059d..7a61673bb 100644 --- a/backend/src/model/stats.rs +++ b/backend/src/model/stats.rs @@ -82,8 +82,12 @@ pub struct SourceStats { } impl SourceStats { - pub fn new(inputs: Vec, targets: Vec)->Self { - Self {inputs, targets} + pub fn new(inputs: Vec, targets: Vec)-> Option { + if inputs.is_empty() || targets.is_empty() { + None + } else { + Some(Self {inputs, targets}) + } } } diff --git a/backend/src/processing/playlist_watch.rs b/backend/src/processing/playlist_watch.rs index 0680c0931..c87091657 100644 --- a/backend/src/processing/playlist_watch.rs +++ b/backend/src/processing/playlist_watch.rs @@ -52,24 +52,26 @@ pub fn process_group_watch(client: &Arc, cfg: &Config, target_n } } +#[derive(Debug, serde::Serialize)] +struct WatchChanges { + pub target: String, + pub group: String, + pub added: Vec, + pub removed: Vec, +} + fn handle_watch_notification(client: &Arc, cfg: &Config, added: &BTreeSet, removed: &BTreeSet, target_name: &str, group_name: &str) { - let added_entries = added.iter().map(std::string::ToString::to_string).collect::>().join("\n\t"); - let removed_entries = removed.iter().map(std::string::ToString::to_string).collect::>().join("\n\t"); + let added = added.iter().map(std::string::ToString::to_string).collect::>(); + let removed = removed.iter().map(std::string::ToString::to_string).collect::>(); + if !added.is_empty() || !removed.is_empty() { + let changes = WatchChanges { + target: target_name.to_string(), + group: group_name.to_string(), + added, + removed + }; - let mut message = vec![]; - if !added_entries.is_empty() { - message.push("added: [\n\t".to_string()); - message.push(added_entries); - message.push("\n]\n".to_string()); - } - if !removed_entries.is_empty() { - message.push("removed: [\n\t".to_string()); - message.push(removed_entries); - message.push("\n]\n".to_string()); - } - - if !message.is_empty() { - let msg = format!("Changes {}/{}\n{}", target_name, group_name, message.join("")); + 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); } diff --git a/backend/src/processing/processor/playlist.rs b/backend/src/processing/processor/playlist.rs index 279c37243..60de26ce9 100644 --- a/backend/src/processing/processor/playlist.rs +++ b/backend/src/processing/processor/playlist.rs @@ -9,7 +9,7 @@ use std::sync::Arc; use std::thread; use tokio::sync::Mutex; -use crate::messaging::send_message; +use crate::messaging::{send_message_json}; use crate::model::Epg; use crate::model::{ConfigTarget, ProcessTargets}; use crate::model::{Mapping}; @@ -376,8 +376,9 @@ 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); - let process_stats = SourceStats::new(input_stats, target_stats); - shared_stats.lock().await.push(process_stats); + if let Some(process_stats) = SourceStats::new(input_stats, target_stats) { + shared_stats.lock().await.push(process_stats); + } }); }, Err(err) => error!("Could not create runtime !!! {err}"), @@ -391,8 +392,9 @@ 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); - let process_stats = SourceStats::new(input_stats, target_stats); - shared_stats.lock().await.push(process_stats); + if let Some(process_stats) = SourceStats::new(input_stats, target_stats) { + shared_stats.lock().await.push(process_stats); + } } drop(update_lock); } @@ -607,7 +609,7 @@ pub async fn exec_processing(client: Arc, app_config: Arc, app_config: Arc, + pub parse_mode: SendMessageParseMode, } fn get_send_message_parse_mode_str(mode: &SendMessageParseMode) -> &'static str { @@ -66,7 +66,7 @@ pub fn telegram_send_message( client: &Arc, instance: &BotInstance, msg: &str, - options: Option, + options: Option<&SendMessageOption>, ) { let chat_id = instance.chat_id.to_string(); let raw_url_str = format!("https://api.telegram.org/bot{}/sendMessage", instance.bot_token); @@ -77,13 +77,13 @@ pub fn telegram_send_message( return; } }; + let request_json_obj = RequestObj { chat_id: instance.chat_id.clone(), message_thread_id: instance.message_thread_id.clone(), text: msg.to_string(), parse_mode: options - .and_then(|o| o.parse_mode) - .map(|mode: SendMessageParseMode| get_send_message_parse_mode_str(&mode)) + .map(|o| get_send_message_parse_mode_str(&o.parse_mode)) .map(ToString::to_string), }; diff --git a/frontend/public/assets/i18n/en.json b/frontend/public/assets/i18n/en.json index 0caec32f5..b012823be 100644 --- a/frontend/public/assets/i18n/en.json +++ b/frontend/public/assets/i18n/en.json @@ -52,6 +52,7 @@ "NOTIFY_ON": "Notify on", "BOT_TOKEN": "Bot Token", "CHAT_IDS": "Chat Ids", + "MARKDOWN": "Markdown", "URL": "Url", "USER": "User", "USERS": "Users", diff --git a/frontend/src/app/components/config/messaging_config_view.rs b/frontend/src/app/components/config/messaging_config_view.rs index dd648a725..0c7ed850c 100644 --- a/frontend/src/app/components/config/messaging_config_view.rs +++ b/frontend/src/app/components/config/messaging_config_view.rs @@ -2,10 +2,7 @@ use crate::app::components::config::config_page::ConfigForm; use crate::app::components::config::config_view_context::ConfigViewContext; use crate::app::components::config::macros::HasFormData; use crate::app::components::{Card, Chip, RadioButtonGroup}; -use crate::{ - config_field, config_field_child, config_field_empty, config_field_hide, config_field_optional, - edit_field_list, edit_field_text, edit_field_text_option, generate_form_reducer, -}; +use crate::{config_field, config_field_bool, config_field_bool_empty, config_field_child, config_field_empty, config_field_hide, config_field_optional, edit_field_bool, edit_field_list, edit_field_text, edit_field_text_option, generate_form_reducer}; use shared::model::{MessagingConfigDto, MsgKind, PushoverMessagingConfigDto, RestMessagingConfigDto, TelegramMessagingConfigDto}; use std::rc::Rc; use std::str::FromStr; @@ -19,6 +16,7 @@ const LABEL_PUSHOVER: &str = "LABEL.PUSHOVER"; const LABEL_REST: &str = "LABEL.REST"; const LABEL_BOT_TOKEN: &str = "LABEL.BOT_TOKEN"; const LABEL_CHAT_IDS: &str = "LABEL.CHAT_IDS"; +const LABEL_MARKDOWN: &str = "LABEL.MARKDOWN"; const LABEL_URL: &str = "LABEL.URL"; const LABEL_TOKEN: &str = "LABEL.TOKEN"; const LABEL_USER: &str = "LABEL.USER"; @@ -29,6 +27,7 @@ generate_form_reducer!( fields { BotToken => bot_token: String, ChatIds => chat_ids: Vec, + Markdown => markdown: bool, } ); @@ -183,6 +182,7 @@ pub fn MessagingConfigView() -> Html { } })} + { config_field_bool!(entry, translate.t(LABEL_MARKDOWN), markdown) } }, None => html! { @@ -190,6 +190,7 @@ pub fn MessagingConfigView() -> Html {

{translate.t(LABEL_TELEGRAM)}

{ config_field_empty!(translate.t(LABEL_BOT_TOKEN)) } { config_field_empty!(translate.t(LABEL_CHAT_IDS)) } + { config_field_bool_empty!(translate.t(LABEL_MARKDOWN)) } }, }; @@ -237,7 +238,7 @@ pub fn MessagingConfigView() -> Html { html! {
{ for notify_on_options.iter().map(|t| { let is_selected = msg_state.form.notify_on.contains(t); - let class = if is_selected { "tp__text-button primary" } else { "tp__text-button" }; + let class = if is_selected { "tp__text-button tp__button-primary" } else { "tp__text-button" }; html! { }}) } @@ -282,8 +283,9 @@ pub fn MessagingConfigView() -> Html {

{translate.t(LABEL_TELEGRAM)}

- { edit_field_text!(telegram_state, translate.t(LABEL_BOT_TOKEN), bot_token, TelegramMessagingConfigFormAction::BotToken) } + { edit_field_text!(telegram_state, translate.t(LABEL_BOT_TOKEN), bot_token, TelegramMessagingConfigFormAction::BotToken, true) } { edit_field_list!(telegram_state, translate.t(LABEL_CHAT_IDS), chat_ids, TelegramMessagingConfigFormAction::ChatIds, translate.t("LABEL.ADD_CHAT_ID")) } + { edit_field_bool!(telegram_state, translate.t(LABEL_MARKDOWN), markdown, TelegramMessagingConfigFormAction::Markdown) }
diff --git a/shared/src/model/config/messaging.rs b/shared/src/model/config/messaging.rs index 00dfa30ce..75746bcd2 100644 --- a/shared/src/model/config/messaging.rs +++ b/shared/src/model/config/messaging.rs @@ -6,6 +6,8 @@ use crate::utils::is_blank_optional_string; pub struct TelegramMessagingConfigDto { pub bot_token: String, pub chat_ids: Vec, + #[serde(default)] + pub markdown: bool, } impl TelegramMessagingConfigDto { diff --git a/shared/src/utils/json_utils.rs b/shared/src/utils/json_utils.rs index 1ca312863..6eeac2818 100644 --- a/shared/src/utils/json_utils.rs +++ b/shared/src/utils/json_utils.rs @@ -3,6 +3,7 @@ use std::io::{self, Read}; use serde::de::DeserializeOwned; use serde::{Deserialize}; use serde_json::{self, Deserializer, Value}; +use crate::utils::humanize_snake_case; fn read_skipping_ws(mut reader: impl Read) -> io::Result { loop { @@ -166,4 +167,52 @@ where { let value: Option = Option::deserialize(deserializer)?; Ok(value.unwrap_or_default()) +} + +const MARKDOWN_SPECIAL_CHARS: &str = r#"_*[]()~`>#+-=|{}.!\"#; + +fn escape_markdown_v2(text: &str) -> String { + let mut escaped = String::new(); + for c in text.chars() { + if MARKDOWN_SPECIAL_CHARS.contains(c) { + escaped.push('\\'); + } + escaped.push(c); + } + escaped +} + +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::Array(arr) => arr.iter() + .map(|v| format!("{pad}\\- {}", format_value(v, indent + 1).trim())) + .collect::>() + .join("\n"), + Value::String(s) => escape_markdown_v2(s), + Value::Number(n) => escape_markdown_v2(&n.to_string()), + Value::Bool(b) => b.to_string(), + Value::Null => "null".to_string(), + } + } + + format_value(value, 0) +} + +pub fn json_str_to_markdown(json_str: &str) -> Result { + let value: Value = serde_json::from_str(json_str)?; + Ok(json_to_markdown(&value)) } \ No newline at end of file diff --git a/shared/src/utils/string_utils.rs b/shared/src/utils/string_utils.rs index 4844bfbcc..4c6f333d4 100644 --- a/shared/src/utils/string_utils.rs +++ b/shared/src/utils/string_utils.rs @@ -106,6 +106,27 @@ pub fn mask_credentials(s: &str) -> String { } } +pub fn humanize_snake_case(s: &str) -> String { + let mut result = String::with_capacity(s.len()); + let mut capitalize_next = true; + + for c in s.chars() { + if c == '_' { + result.push(' '); + capitalize_next = true; + } else if capitalize_next { + for up in c.to_uppercase() { + result.push(up); + } + capitalize_next = false; + } else { + result.push(c); + } + } + + result +} + #[cfg(test)] mod test { use std::collections::HashSet; From 79906952c5676eaebbc7fc1f98ac09cb0b34d705 Mon Sep 17 00:00:00 2001 From: euzu Date: Wed, 22 Oct 2025 11:33:49 +0200 Subject: [PATCH 3/3] Telegram markdown support for structured json messages --- README.md | 2 +- backend/src/messaging.rs | 20 +++--- backend/src/model/stats.rs | 2 +- backend/src/processing/playlist_watch.rs | 2 +- backend/src/processing/processor/playlist.rs | 74 +++++++++++--------- backend/src/utils/network/xtream.rs | 6 +- shared/src/utils/json_utils.rs | 33 +++++---- 7 files changed, 77 insertions(+), 62 deletions(-) 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),