diff --git a/CHANGELOG.md b/CHANGELOG.md index cc63d1d36..42746a8c0 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -6,6 +6,7 @@ * Added Xtream Api Endpoints. * Added multiple input support * Added Telegram message support +* Added Target watch for groups # v1.0.1(2023-09-07) * Refactored sorting. Sorting channels inside group now possible diff --git a/Cargo.lock b/Cargo.lock index 45f870b74..cb9c0ad03 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -378,6 +378,15 @@ version = "0.21.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9ba43ea6f343b788c8764558649e08df62f86c6ef251fdaeb1ffd010a9ae50a2" +[[package]] +name = "bincode" +version = "1.3.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b1f45e9417d87227c7a56d22e471c6206462cba514c7590c09aff4cf6d1ddcad" +dependencies = [ + "serde", +] + [[package]] name = "bitflags" version = "1.3.2" @@ -1275,6 +1284,7 @@ dependencies = [ "actix-rt", "actix-server", "actix-web", + "bincode", "chrono", "clap", "cron", diff --git a/Cargo.toml b/Cargo.toml index a84ba1413..19e47adb4 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -38,3 +38,4 @@ mime = "0.3.17" log = "0.4.20" env_logger = "0.10.0" rustelebot = "0.3.2" +bincode = "1.3.3" diff --git a/README.md b/README.md index c6221c21c..8b8c490d9 100644 --- a/README.md +++ b/README.md @@ -163,6 +163,7 @@ Has the following top level entries: * `filter` _mandatory_, * `rename` _optional_ * `mapping` _optional_ +* `watch` _optional_ ### 1.5.2.1 `publish` and `filename` @@ -314,6 +315,21 @@ sources: new_name: 1. DE$1 ``` +### 1.5.2.9 `watch` +For each target with a *unique name*, you can define a watched groups. +It is a list of final group names from this target playlist. +Final means in this case: the name in the resulting playlist after applying all steps +of transformation. + +For example given the following configuration: +```yaml +watch: + - 'FR | Movies' + - 'FR | Series' +``` + +Changes from this groups will be printed as info on console and send to +the configured messaging (f.e. telegram channel). ### 1.6 `messaging` `messaging` is an optional configuration for receiving messages. diff --git a/config.yml b/config.yml index 194643d4a..677dd9d52 100644 --- a/config.yml +++ b/config.yml @@ -48,6 +48,8 @@ sources: - Belgique - Suisse - Canada + watch: + - 'FR - Movies' - filename: playlist_strm type: strm filter: 'Group ~ "^tv-shows.*"' diff --git a/src/download.rs b/src/download.rs index cc465f0ae..b0a31baae 100644 --- a/src/download.rs +++ b/src/download.rs @@ -1,10 +1,11 @@ use std::path::PathBuf; use std::sync::atomic::AtomicI32; use log::debug; -use crate::{m3u_parser, utils, xtream_parser}; +use crate::{utils}; use crate::m3u_filter_error::M3uFilterError; use crate::model::config::{Config, ConfigInput}; use crate::model::model_m3u::{PlaylistGroup, XtreamCluster}; +use crate::processing::{m3u_parser, xtream_parser}; fn prepare_file_path(input: &ConfigInput, working_dir: &String, action: &str) -> Option { let persist_file: Option = diff --git a/src/main.rs b/src/main.rs index b43bd7c1f..cb6c70f92 100644 --- a/src/main.rs +++ b/src/main.rs @@ -7,15 +7,13 @@ mod m3u_filter_error; mod config_reader; mod model; mod filter; -mod m3u_parser; -mod playlist_processor; mod repository; mod download; mod utils; mod messaging; -mod xtream_parser; mod test; mod api; +mod processing; use env_logger::{Builder}; use log::{debug, error, info, LevelFilter}; @@ -25,6 +23,7 @@ use crate::config_reader::{read_api_proxy_config, read_config, read_mappings}; use crate::messaging::send_message; use crate::model::config::{Config, ProcessTargets, validate_targets}; use crate::m3u_filter_error::{M3uFilterErrorKind}; +use crate::processing::playlist_processor; #[derive(Parser)] #[command(name = "m3u-filter")] diff --git a/src/model/config.rs b/src/model/config.rs index 9b29630be..c1052c887 100644 --- a/src/model/config.rs +++ b/src/model/config.rs @@ -14,7 +14,7 @@ use crate::utils::get_working_path; fn default_as_frm() -> ProcessingOrder { ProcessingOrder::Frm } -fn default_as_default() -> String { String::from("default") } +pub(crate) fn default_as_default() -> String { String::from("default") } fn default_as_empty_map() -> HashMap { HashMap::new() } @@ -166,6 +166,7 @@ pub(crate) struct ConfigTarget { pub mapping: Option>, #[serde(default = "default_as_frm")] pub processing_order: ProcessingOrder, + pub watch: Option>, #[serde(skip_serializing, skip_deserializing)] pub _filter: Option, #[serde(skip_serializing, skip_deserializing)] diff --git a/src/m3u_parser.rs b/src/processing/m3u_parser.rs similarity index 100% rename from src/m3u_parser.rs rename to src/processing/m3u_parser.rs diff --git a/src/processing/mod.rs b/src/processing/mod.rs new file mode 100644 index 000000000..14c4791c8 --- /dev/null +++ b/src/processing/mod.rs @@ -0,0 +1,4 @@ +pub(crate) mod m3u_parser; +pub(crate) mod xtream_parser; +pub(crate) mod playlist_processor; +pub(crate) mod playlist_watch; \ No newline at end of file diff --git a/src/playlist_processor.rs b/src/processing/playlist_processor.rs similarity index 96% rename from src/playlist_processor.rs rename to src/processing/playlist_processor.rs index 2919c41c9..5d3b72a93 100644 --- a/src/playlist_processor.rs +++ b/src/processing/playlist_processor.rs @@ -6,7 +6,7 @@ use std::thread; use log::{debug, error, info}; use unidecode::unidecode; use crate::{model::config, Config, valid_property, create_m3u_filter_error_result}; -use crate::model::config::{ConfigTarget, InputAffix, InputType, ProcessTargets}; +use crate::model::config::{ConfigTarget, default_as_default, InputAffix, InputType, ProcessTargets}; use crate::model::model_config::{SortOrder::{Asc, Desc}, ItemField, AFFIX_FIELDS, ProcessingOrder}; use crate::filter::{get_field_value, MockValueProcessor, set_field_value, ValueProvider}; use crate::repository::m3u_repository::write_playlist; @@ -14,6 +14,7 @@ use crate::model::model_m3u::{FetchedPlaylist, FieldAccessor, PlaylistGroup, Pla use crate::model::mapping::{Mapping, MappingValueProcessor}; use crate::download::{get_m3u_playlist, get_xtream_playlist}; use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; +use crate::processing::playlist_watch::process_group_watch; use crate::repository::xtream_repository::xtream_save_playlist; fn filter_playlist(playlist: &mut [PlaylistGroup], target: &ConfigTarget) -> Option> { @@ -385,6 +386,20 @@ pub(crate) fn process_playlist(playlists: &mut [FetchedPlaylist], if !new_playlist.is_empty() { sort_playlist(target, &mut new_playlist); + + if target.watch.is_some() { + if default_as_default().eq_ignore_ascii_case(&target.name) { + error!("cant watch a target with no unique name"); + } else { + let titles = target.watch.as_ref().unwrap(); + new_playlist.iter().for_each(|pl| { + if titles.contains(&pl.title) { + process_group_watch(cfg, &target.name, pl) + } + }); + } + } + let publish = target.publish; if target.filename.is_some() { let result = write_playlist(target, cfg, &mut new_playlist); @@ -408,4 +423,3 @@ pub(crate) fn process_playlist(playlists: &mut [FetchedPlaylist], Ok(()) } } - diff --git a/src/processing/playlist_watch.rs b/src/processing/playlist_watch.rs new file mode 100644 index 000000000..d2a3e8936 --- /dev/null +++ b/src/processing/playlist_watch.rs @@ -0,0 +1,87 @@ +use std::collections::BTreeSet; +use std::path::PathBuf; +use log::{error, info}; +use regex::Regex; +use crate::messaging::send_message; +use crate::model::config::Config; +use crate::model::model_m3u::PlaylistGroup; +use crate::utils; + +pub(crate) fn process_group_watch(cfg: &Config, target_name: &str, pl: &PlaylistGroup) { + let mut new_tree = BTreeSet::new(); + pl.channels.iter().for_each(|chan| { + let header = chan.header.borrow(); + let title = if header.title.is_empty() { header.title.to_string() } else { header.name.to_string() }; + new_tree.insert(title); + }); + + let filename_re = Regex::new(r"[^A-Za-z0-9_.-\\s]").unwrap(); + let watch_filename = format!("watch_{}.bin", filename_re.replace_all(target_name, "_")); + match utils::get_file_path(&cfg.working_dir, Some(std::path::PathBuf::from(&watch_filename))) { + Some(path) => { + let save_path = path.clone(); + if path.exists() { + match load_tree(&path) { + Some(loaded_tree) => { + // Find elements in set2 but not in set1 + let added_difference: BTreeSet = new_tree.difference(&loaded_tree).cloned().collect(); + let removed_difference: BTreeSet = loaded_tree.difference(&new_tree).cloned().collect(); + if !added_difference.is_empty() || !removed_difference.is_empty() { + handle_watch_notification(cfg, added_difference, removed_difference); + } + } + None => { + error!("failed to load watch_file {}", &path.into_os_string().into_string().unwrap()); + } + } + } + match save_tree(&save_path, new_tree) { + Ok(_) => {} + Err(err) => { + error!("failed to write watch_file {}: {}", &save_path.into_os_string().into_string().unwrap(), err) + } + } + } + None => { + error!("failed to write watch_file {}", &watch_filename); + } + } +} + +fn handle_watch_notification(cfg: &Config, added: BTreeSet, removed: BTreeSet) { + let added_entries = added.iter().map(|name| name.to_string()).collect::>().join("\n\t"); + let removed_entries = removed.iter().map(|name| name.to_string()).collect::>().join("\n\t"); + + 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() { + info!("{}", message.join("").as_str()); + send_message(&cfg.messaging, message.join("").as_str()) + } +} + +fn load_tree(path: &PathBuf) -> Option> { + match std::fs::read(path) { + Ok(encoded) => { + let decoded: BTreeSet = bincode::deserialize(&encoded[..]).unwrap(); + Some(decoded) + } + Err(_) => None, + } +} + +fn save_tree(path: &PathBuf, tree: BTreeSet) -> std::io::Result<()> { + let encoded: Vec = bincode::serialize(&tree).unwrap(); + std::fs::write(path, encoded) +} + diff --git a/src/xtream_parser.rs b/src/processing/xtream_parser.rs similarity index 100% rename from src/xtream_parser.rs rename to src/processing/xtream_parser.rs