diff --git a/CHANGELOG.md b/CHANGELOG.md index 68b13b765..dd45a86b0 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -21,6 +21,7 @@ schedules: - series_channels - schedule: "0 0 20 * * * *" ``` +- Stats have now target information # 2.0.10 (2024-12-03) - added Target Output Option `m3u_include_type_in_url`, default false. This adds `live`, `movie`, `series` to the url of the stream in reverse proxy mode. diff --git a/src/model/stats.rs b/src/model/stats.rs index 0709cd874..07504a2b3 100644 --- a/src/model/stats.rs +++ b/src/model/stats.rs @@ -1,4 +1,4 @@ -use std::fmt::Display; +use std::fmt::{Display}; use serde::{Serialize, Serializer}; use crate::model::config::InputType; @@ -48,4 +48,48 @@ impl Display for InputStats { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { serde_json::to_string(&self).map_or(Err(std::fmt::Error), |json_str| write!(f, "{json_str}")) } -} \ No newline at end of file +} + + +#[derive(Debug, Clone, Serialize)] +pub struct TargetStats { + #[serde(rename = "target")] + pub name: String, + pub success: bool, +} + +impl TargetStats { + pub fn success(name: &str) -> Self { + Self {name: name.to_string(), success: true} + } + pub fn failure(name: &str) -> Self { + Self {name: name.to_string(), success: false} + } +} + +impl Display for TargetStats { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + serde_json::to_string(&self).map_or(Err(std::fmt::Error), |json_str| write!(f, "{json_str}")) + } +} + +#[derive(Debug, Clone, Serialize)] +pub struct SourceStats { + #[serde(rename = "inputs")] + pub inputs: Vec, + #[serde(rename = "targets")] + pub targets: Vec, +} + +impl SourceStats { + pub fn new(inputs: Vec, targets: Vec)->Self { + Self {inputs, targets} + } +} + +impl Display for SourceStats { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + serde_json::to_string(&self).map_or(Err(std::fmt::Error), |json_str| write!(f, "{json_str}")) + } +} + diff --git a/src/processing/playlist_processor.rs b/src/processing/playlist_processor.rs index e9c6a8729..18d697626 100644 --- a/src/processing/playlist_processor.rs +++ b/src/processing/playlist_processor.rs @@ -20,7 +20,7 @@ use crate::model::config::{ConfigSortChannel, ConfigSortGroup, ConfigTarget, Inp ItemField, ProcessTargets, ProcessingOrder, SortOrder::{Asc, Desc}}; use crate::model::mapping::{CounterModifier, Mapping, MappingValueProcessor}; use crate::model::playlist::{FetchedPlaylist, FieldGetAccessor, FieldSetAccessor, PlaylistGroup, PlaylistItem, XtreamCluster}; -use crate::model::stats::{InputStats, PlaylistStats}; +use crate::model::stats::{SourceStats, InputStats, TargetStats, PlaylistStats}; use crate::processing::affix_processor::apply_affixes; use crate::processing::playlist_watch::process_group_watch; use crate::processing::xmltv_parser::flatten_tvguide; @@ -289,10 +289,11 @@ fn is_target_enabled(target: &ConfigTarget, user_targets: &ProcessTargets) -> bo (!user_targets.enabled && target.enabled) || (user_targets.enabled && user_targets.has_target(target.id)) } -async fn process_source(cfg: Arc, source_idx: usize, user_targets: Arc) -> (Vec, Vec) { +async fn process_source(cfg: Arc, source_idx: usize, user_targets: Arc) -> (Vec, Vec, Vec) { let source = cfg.sources.get(source_idx).unwrap(); let mut errors = vec![]; - let mut stats = HashMap::::new(); + let mut input_stats = HashMap::::new(); + let mut target_stats = Vec::::new(); let mut source_playlists = Vec::new(); let enabled_inputs = source.inputs.iter().filter(|item| item.enabled).count(); // Downlod the sources @@ -330,8 +331,8 @@ async fn process_source(cfg: Arc, source_idx: usize, user_targets: Arc

, source_idx: usize, user_targets: Arc

()); for target in &source.targets { if is_target_enabled(target, &user_targets) { - match process_playlist(&mut source_playlists, target, &cfg, &mut stats, &mut errors).await { - Ok(()) => {} - Err(mut err) => errors.append(&mut err) + match process_playlist(&mut source_playlists, target, &cfg, &mut input_stats, &mut errors).await { + Ok(()) => { + target_stats.push(TargetStats::success(&target.name)); + } + Err(mut err) => { + target_stats.push(TargetStats::failure(&target.name)); + errors.append(&mut err); + } } } } } - (stats.into_values().collect(), errors) + (input_stats.into_values().collect(), target_stats, errors) } fn create_input_stat(group_count: usize, channel_count: usize, error_count: usize, input_type: InputType, input_name: &str, secs_took: u64) -> InputStats { @@ -368,7 +374,7 @@ fn create_input_stat(group_count: usize, channel_count: usize, error_count: usiz } } -async fn process_sources(config: Arc, user_targets: Arc) -> (Vec, Vec) { +async fn process_sources(config: Arc, user_targets: Arc) -> (Vec, Vec) { let mut handle_list = vec![]; let thread_num = config.threads; let process_parallel = thread_num > 1 && config.sources.len() > 1; @@ -376,7 +382,7 @@ async fn process_sources(config: Arc, user_targets: Arc) debug!("Using {} threads", thread_num); } let errors = Arc::new(Mutex::>::new(vec![])); - let stats = Arc::new(Mutex::>::new(vec![])); + let stats = Arc::new(Mutex::>::new(vec![])); for (index, _) in config.sources.iter().enumerate() { let shared_errors = errors.clone(); let shared_stats = stats.clone(); @@ -386,9 +392,10 @@ async fn process_sources(config: Arc, user_targets: Arc) let handles = &mut handle_list; let process = move || { System::new().block_on(async { - let (mut res_stats, mut res_errors) = process_source(cfg, index, usr_trgts).await; + let (input_stats, target_stats, mut res_errors) = process_source(cfg, index, usr_trgts).await; shared_errors.lock().await.append(&mut res_errors); - shared_stats.lock().await.append(&mut res_stats); + let process_stats = SourceStats::new(input_stats, target_stats); + shared_stats.lock().await.push(process_stats); }); }; handles.push(thread::spawn(process)); @@ -396,9 +403,10 @@ async fn process_sources(config: Arc, user_targets: Arc) handles.drain(..).for_each(|handle| { let _ = handle.join(); }); } } else { - let (mut res_stats, mut res_errors) = process_source(cfg, index, usr_trgts).await; + let (input_stats, target_stats, mut res_errors) = process_source(cfg, index, usr_trgts).await; shared_errors.lock().await.append(&mut res_errors); - shared_stats.lock().await.append(&mut res_stats); + let process_stats = SourceStats::new(input_stats, target_stats); + shared_stats.lock().await.push(process_stats); } } for handle in handle_list {