replaced async_std

This commit is contained in:
euzu
2025-02-06 10:49:34 +01:00
parent bbf6d6e014
commit 75c63ae237
35 changed files with 486 additions and 715 deletions
+6 -6
View File
@@ -5,7 +5,7 @@ use crate::model::config::ConfigRename;
use crate::utils::network::epg;
use crate::utils::network::m3u;
use crate::utils::network::xtream;
use async_std::sync::Mutex;
use parking_lot::Mutex;
use core::cmp::Ordering;
use std::cell::RefCell;
use std::collections::{HashMap, HashSet};
@@ -399,7 +399,7 @@ async fn process_sources(client: Arc<reqwest::Client>, config: Arc<Config>, user
for (index, _) in config.sources.iter().enumerate() {
// We're using the file lock this way on purpose
let source_lock_path = PathBuf::from(format!("source_{index}"));
let Ok(update_lock) = config.file_locks.try_write_lock(&source_lock_path).await else {
let Ok(update_lock) = config.file_locks.try_write_lock(&source_lock_path) else {
warn!("The update operation for the source at index {index} was skipped because an update is already in progress.");
continue;
};
@@ -414,9 +414,9 @@ async fn process_sources(client: Arc<reqwest::Client>, config: Arc<Config>, user
let process = move || {
System::new().block_on(async {
let (input_stats, target_stats, mut res_errors) = process_source(Arc::clone(&http_client), cfg, index, usr_trgts).await;
shared_errors.lock().await.append(&mut res_errors);
shared_errors.lock().append(&mut res_errors);
let process_stats = SourceStats::new(input_stats, target_stats);
shared_stats.lock().await.push(process_stats);
shared_stats.lock().push(process_stats);
});
};
handles.push(thread::spawn(process));
@@ -425,9 +425,9 @@ async fn process_sources(client: Arc<reqwest::Client>, config: Arc<Config>, user
}
} else {
let (input_stats, target_stats, mut res_errors) = process_source(Arc::clone(&client), cfg, index, usr_trgts).await;
shared_errors.lock().await.append(&mut res_errors);
shared_errors.lock().append(&mut res_errors);
let process_stats = SourceStats::new(input_stats, target_stats);
shared_stats.lock().await.push(process_stats);
shared_stats.lock().push(process_stats);
}
drop(update_lock);
}
+8 -10
View File
@@ -1,8 +1,8 @@
use crate::m3u_filter_error::{info_err, notify_err};
use crate::m3u_filter_error::{str_to_io_error, to_io_error, M3uFilterError, M3uFilterErrorKind};
use crate::model::config::{Config, ConfigInput};
use crate::model::playlist::{FetchedPlaylist, PlaylistEntry, PlaylistItem, PlaylistItemType, XtreamCluster};
use crate::repository::storage::get_input_storage_path;
use crate::m3u_filter_error::{info_err, notify_err};
use serde::{Deserialize, Serialize};
use std::collections::HashMap;
use std::fs::File;
@@ -107,16 +107,14 @@ where
}
};
match cfg.file_locks.read_lock(&file_path).await {
Ok(file_lock) => {
if let Ok(info_records) = BPlusTree::<u32, V>::load(&file_path) {
info_records.iter().for_each(|(provider_id, record)| {
processed_info_ids.insert(*provider_id, extract_ts(record));
});
}
drop(file_lock);
{
let file_lock = cfg.file_locks.read_lock(&file_path);
if let Ok(info_records) = BPlusTree::<u32, V>::load(&file_path) {
info_records.iter().for_each(|(provider_id, record)| {
processed_info_ids.insert(*provider_id, extract_ts(record));
});
}
Err(err) => errors.push(info_err!(format!("{err}"))),
drop(file_lock);
}
processed_info_ids
}
+1 -4
View File
@@ -139,10 +139,7 @@ async fn process_series_info(
return result;
};
let Ok(_file_lock) = cfg.file_locks.read_lock(&info_path).await else {
errors.push(notify_err!("Could not lock input info file for series".to_string()));
return result;
};
let _file_lock = cfg.file_locks.read_lock(&info_path);
// Contains the Series Info with episode listing
let Ok(mut info_reader) = IndexedDocumentReader::<u32, String>::new(&info_path, &idx_path) else { return result; };