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;