resolve series wip

This commit is contained in:
euzu
2024-12-25 14:12:30 +01:00
parent 454e85e6d2
commit dfe6269866
4 changed files with 50 additions and 59 deletions
+12 -10
View File
@@ -11,7 +11,6 @@ use std::thread;
use actix_rt::System;
use log::{debug, error, info, log_enabled, trace, Level};
use std::time::Instant;
use unidecode::unidecode;
use crate::filter::{get_field_value, set_field_value, MockValueProcessor, ValueProvider};
@@ -531,19 +530,22 @@ fn process_watch(target: &ConfigTarget, cfg: &Config, new_playlist: &Vec<Playlis
}
pub async fn exec_processing(cfg: Arc<Config>, targets: Arc<ProcessTargets>) {
let start_time = Instant::now();
let (stats, errors) = process_sources(cfg.clone(), targets.clone()).await;
// log errors
for err in &errors {
error!("{}", err.message);
errors.iter().for_each(|err| error!("{}", err.message));
if let Ok(stats_msg) = serde_json::to_string(&serde_json::Value::Object(serde_json::map::Map::from_iter([("stats".to_string(), serde_json::to_value(stats).unwrap())]))) {
// print stats
info!("{}", stats_msg);
// send stats
send_message(&MsgKind::Stats, cfg.messaging.as_ref(), stats_msg.as_str());
}
let stats_msg = format!("{{\"stats\": {}}}", stats.iter().map(std::string::ToString::to_string).collect::<Vec<String>>().join("\n"));
// print stats
info!("{}", stats_msg);
// send stats
send_message(&MsgKind::Stats, cfg.messaging.as_ref(), stats_msg.as_str());
// send errors
if let Some(message) = get_errors_notify_message!(errors, 255) {
let error_msg = format!("{{\"errors\": \"{}\"}}", message.as_str());
send_message(&MsgKind::Error, cfg.messaging.as_ref(), error_msg.as_str());
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(&MsgKind::Error, cfg.messaging.as_ref(), error_msg.as_str());
}
}
let elapsed = start_time.elapsed().as_secs();
info!("Update process finished! Took {elapsed} secs.");
}
+16 -26
View File
@@ -7,7 +7,7 @@ use crate::{info_err, notify_err};
use serde::{Deserialize, Serialize};
use serde_json::Value;
use std::collections::HashMap;
use std::fs::{File};
use std::fs::File;
use std::io::{BufWriter, Error, ErrorKind, Write};
use std::path::PathBuf;
@@ -130,32 +130,22 @@ pub(in crate::processing) fn create_resolve_info_wal_files(cfg: &Config, input:
}
}
pub(in crate::processing) fn has_different_ts(ts: u64, pli: &PlaylistItem, field: &str) -> bool {
pli.header
.borrow()
.additional_properties
.as_ref()
.map_or(false, |v| match v {
Value::Object(map) => {
if let Some(updated) = map.get(field) {
if let Some(update_ts) = get_u64_from_serde_value(updated) {
return update_ts != ts;
}
}
true
}
_ => true,
})
}
pub(in crate::processing) fn should_update_info(pli: &PlaylistItem, processed_provider_ids: &HashMap<u32, u64>, field: &str) -> (bool, u32, u64) {
if let Some(provider_id) = pli.header.borrow_mut().get_provider_id() {
let timestamp = processed_provider_ids.get(&provider_id);
(timestamp.is_none() || has_different_ts(*timestamp.unwrap(), pli, field), provider_id, *timestamp.unwrap_or(&0))
} else {
(false, 0, 0)
}
let Some(provider_id) = pli.header.borrow_mut().get_provider_id() else { return (false, 0, 0) };
let last_modified = pli.header.borrow().additional_properties.as_ref().map_or(None, |v| match v {
Value::Object(map) => {
if let Some(updated) = map.get(field) {
get_u64_from_serde_value(updated)
} else {
None
}
}
_ => None,
});
let old_timestamp = processed_provider_ids.get(&provider_id);
(old_timestamp.is_none()
|| last_modified.is_none()
|| *old_timestamp.unwrap() != last_modified.unwrap(), provider_id, last_modified.unwrap_or(0))
}
pub(in crate::processing) async fn read_processed_info_ids<V, F>(cfg: &Config, errors: &mut Vec<M3uFilterError>, fpl: &FetchedPlaylist<'_>,