new messaging option

This commit is contained in:
euzu
2024-01-11 09:43:16 +01:00
parent 2b511bba66
commit d16fd8a6d4
12 changed files with 85 additions and 14 deletions
+1
View File
@@ -213,6 +213,7 @@ async fn xtream_get_stream_info(app_state: &AppState, target_name: &str, stream_
if response.status().is_success() {
match response.text().await {
Ok(content) => {
// TODO we are not replacing direct_source, we should add an option to do this.
xtream_persist_stream_info(app_state, target_name, stream_id, cluster,
target_input, content.as_str()).await;
return Ok(content);
+19
View File
@@ -1,4 +1,5 @@
use log::{debug, error};
use reqwest::header;
use crate::model::config::{MessagingConfig};
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, PartialEq)]
@@ -9,6 +10,8 @@ pub(crate) enum MsgKind {
Stats,
#[serde(rename = "error")]
Error,
#[serde(rename = "watch")]
Watch,
}
fn is_enabled(kind: &MsgKind, cfg: &MessagingConfig) -> bool {
@@ -28,6 +31,22 @@ pub(crate) fn send_message(kind: &MsgKind, cfg: &Option<MessagingConfig>, msg: &
}
};
}
if let Some(rest) = &messaging.rest {
let url = rest.url.to_owned();
let data = msg.to_owned();
actix_rt::spawn(async move {
let client = reqwest::Client::new();
match client.post(&url)
.header(header::CONTENT_TYPE, mime::APPLICATION_JSON.to_string())
.body(data)
.send()
.await {
Ok(_) => debug!("Text message sent successfully to rest api"),
Err(e) => error!("Text message wasn't sent to rest api because of: {}", e)
}
});
}
}
}
}
+6
View File
@@ -445,11 +445,17 @@ pub(crate) struct TelegramMessagingConfig {
pub chat_ids: Vec<String>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub(crate) struct RestMessagingConfig {
pub url: String,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub(crate) struct MessagingConfig {
#[serde(default = "default_as_empty_list")]
pub notify_on: Vec<MsgKind>,
pub telegram: Option<TelegramMessagingConfig>,
pub rest: Option<RestMessagingConfig>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
+2 -2
View File
@@ -8,7 +8,7 @@ pub(crate) struct PlaylistStats {
impl ToString for PlaylistStats {
fn to_string(&self) -> String {
format!("{{groups: {}, channels: {}}}", self.group_count, self.channel_count)
format!("{{\"groups\": {}, \"channels\": {}}}", self.group_count, self.channel_count)
}
}
@@ -23,7 +23,7 @@ pub(crate) struct InputStats {
impl ToString for InputStats {
fn to_string(&self) -> String {
format!("{{name: {}, type: {}, errors: {}, raw: {}, processed: {}}}",
format!("{{\"name\": {}, \"type\": {}, \"errors\": {}, \"raw\": {}, \"processed\": {}}}",
self.name, self.input_type.to_string(), self.error_count,
self.raw_stats.to_string(), self.processed_stats.to_string())
}
+3 -2
View File
@@ -535,7 +535,7 @@ fn persist_playlist(playlist: &[PlaylistGroup], epg: Option<Epg>,
pub(crate) async fn exec_processing(cfg: Arc<Config>, targets: Arc<ProcessTargets>) {
let (stats, errors) = process_sources(cfg.to_owned(), targets.to_owned()).await;
let stats_msg = format!("Stats: {}", stats.iter().map(|stat| stat.to_string()).collect::<Vec<String>>().join("\n"));
let stats_msg = format!("{{\"stats\": {}}}", stats.iter().map(|stat| stat.to_string()).collect::<Vec<String>>().join("\n"));
// print stats
info!("{}", stats_msg);
// send stats
@@ -544,6 +544,7 @@ pub(crate) async fn exec_processing(cfg: Arc<Config>, targets: Arc<ProcessTarget
errors.iter().for_each(|err| error!("{}", err.message));
// send errors
if let Some(message) = get_errors_notify_message!(errors, 255) {
send_message(&MsgKind::Error, &cfg.messaging, message.as_str());
let error_msg = format!("{{\"errors\": \"{}\"}}",message.as_str());
send_message(&MsgKind::Error, &cfg.messaging, error_msg.as_str());
}
}
+1 -1
View File
@@ -72,7 +72,7 @@ fn handle_watch_notification(cfg: &Config, added: BTreeSet<String>, removed: BTr
if !message.is_empty() {
let msg = format!("Changes {}/{}\n{}", target_name, group_name, message.join(""));
info!("{}", &msg);
send_message(&MsgKind::Info, &cfg.messaging, &msg);
send_message(&MsgKind::Watch, &cfg.messaging, &msg);
}
}