diff --git a/src/api/m3u_api.rs b/src/api/m3u_api.rs index 13fcda6ce..fab63837f 100644 --- a/src/api/m3u_api.rs +++ b/src/api/m3u_api.rs @@ -41,7 +41,7 @@ async fn m3u_api_stream( match get_target_storage_path(&app_state.config, target.name.as_str()) { Some(target_path) => { if let Some((m3u_path, idx_path)) = m3u_get_file_paths(&target_path) { - match m3u_get_item_for_stream_id(m3u_stream_id, &m3u_path, &idx_path) { + match m3u_get_item_for_stream_id(&app_state.config, m3u_stream_id, &m3u_path, &idx_path) { Ok(m3u_item) => { return stream_response(m3u_item.url.as_str(), &req, None).await; } diff --git a/src/api/scheduler.rs b/src/api/scheduler.rs index 4c0d0c6e1..a7286c739 100644 --- a/src/api/scheduler.rs +++ b/src/api/scheduler.rs @@ -1,29 +1,81 @@ use std::str::FromStr; -use std::time::Duration; +use std::time::{Duration, Instant, SystemTime}; use actix_web::web::Data; -use chrono::Local; +use chrono::{DateTime, FixedOffset, Local}; use cron::Schedule; use log::error; use crate::api::api_model::AppState; use crate::exit; use crate::processing::playlist_processor::exec_processing; +fn datetime_to_instant(datetime: DateTime) -> Instant { + // Convert DateTime to SystemTime + let target_system_time: SystemTime = datetime.into(); + + // Get the current SystemTime + let now_system_time = SystemTime::now(); + + // Calculate the duration between now and the target time + let duration_until = target_system_time + .duration_since(now_system_time) + .unwrap_or_else(|_| Duration::from_secs(0)); + + // Get the current Instant and add the duration to calculate the target Instant + Instant::now() + duration_until +} + pub(crate) async fn start_scheduler(expression: &str, data: Data) -> ! { match Schedule::from_str(expression) { Ok(schedule) => { let offset = *Local::now().offset(); loop { let mut upcoming = schedule.upcoming(offset).take(1); - actix_rt::time::sleep(Duration::from_millis(500)).await; - let local = &Local::now(); - if let Some(datetime) = upcoming.next() { - if datetime.timestamp() <= local.timestamp() { - exec_processing(data.config.clone(), data.targets.clone()).await; - } + actix_rt::time::sleep_until(actix_rt::time::Instant::from(datetime_to_instant(datetime))).await; + exec_processing(data.config.clone(), data.targets.clone()).await; } } } Err(err) => exit!("Failed to start scheduler: {}", err) } +} + +#[cfg(test)] +mod tests { + use std::str::FromStr; + use std::sync::atomic::{AtomicU8, Ordering}; + use chrono::Local; + use cron::Schedule; + use crate::api::scheduler::datetime_to_instant; + + #[actix_rt::test] + async fn test_run_scheduler() { + // Define a cron expression that runs every second + let expression = "0/1 * * * * * *"; // every second + + let runs = AtomicU8::new(0); + let run_me = || runs.fetch_add(1, Ordering::Relaxed); + + let start = std::time::Instant::now(); + match Schedule::from_str(expression) { + Ok(schedule) => { + let offset = *Local::now().offset(); + loop { + let mut upcoming = schedule.upcoming(offset).take(1); + if let Some(datetime) = upcoming.next() { + actix_rt::time::sleep_until(actix_rt::time::Instant::from(datetime_to_instant(datetime))).await; + run_me(); + } + if runs.load(Ordering::Relaxed) == 6 { + break; + } + } + } + Err(_) => {} + }; + let duration = start.elapsed(); + + assert!(runs.load(Ordering::Relaxed) == 6, "Failed to run"); + assert!(duration.as_secs() > 4, "Failed time"); + } } \ No newline at end of file diff --git a/src/api/xtream_api.rs b/src/api/xtream_api.rs index fa05a3c30..e095affaf 100644 --- a/src/api/xtream_api.rs +++ b/src/api/xtream_api.rs @@ -16,11 +16,49 @@ use crate::api::api_utils::{get_user_server_info, get_user_target, get_user_targ use crate::model::api_proxy::{ProxyType, ProxyUserCredentials}; use crate::model::config::{Config, ConfigInput, ConfigTarget}; use crate::model::config::TargetType; -use crate::model::playlist::{XtreamCluster, XtreamPlaylistItem}; +use crate::model::playlist::{PlaylistItemType, XtreamCluster, XtreamPlaylistItem}; use crate::model::xtream::XtreamMappingOptions; +use crate::repository::storage::{get_target_storage_path, hash_string}; +use crate::repository::target_id_mapping::TargetIdMapping; use crate::repository::xtream_repository; use crate::utils::{json_utils, request_utils}; + +macro_rules! try_option_bad_request { + ($option:expr, $msg_is_error:expr, $msg:expr) => { + match $option { + Some(value) => value, + None => { + if $msg_is_error {error!("{}", $msg);} else {debug!("{}", $msg);} + return HttpResponse::BadRequest().finish(); + } + } + }; + ($option:expr) => { + match $option { + Some(value) => value, + None => return HttpResponse::BadRequest().finish(), + } + }; +} +macro_rules! try_result_bad_request { + ($option:expr, $msg_is_error:expr, $msg:expr) => { + match $option { + Ok(value) => value, + Err(_) => { + if $msg_is_error {error!("{}", $msg);} else {debug!("{}", $msg);} + return HttpResponse::BadRequest().finish(); + } + } + }; + ($option:expr) => { + match $option { + Ok(value) => value, + Err(_) => return HttpResponse::BadRequest().finish(), + } + }; +} + enum XtreamApiStreamContext { LiveAlt, Live, @@ -107,7 +145,6 @@ fn get_xtream_player_api_stream_url(input: &ConfigInput, context: &str, action_p } } - fn get_user_info(user: &ProxyUserCredentials, cfg: &Config) -> XtreamAuthorizationResponse { let server_info = get_user_server_info(cfg, user); @@ -155,43 +192,34 @@ async fn xtream_player_api_stream( app_state: &web::Data, stream_req: XtreamApiStreamRequest<'_>, ) -> HttpResponse { - if let Some((user, target)) = get_user_target_by_credentials(stream_req.username, stream_req.password, api_req, app_state) { - let target_name = &target.name; - if target.has_output(&TargetType::Xtream) { - let (action_stream_id, stream_ext) = xtream_api_request_separate_number_and_rest(stream_req.stream_id); - let virtual_id: u32 = match FromStr::from_str(action_stream_id.trim()) { - Ok(id) => id, - Err(_) => return HttpResponse::BadRequest().finish() - }; - - if let Ok(pli) = xtream_repository::xtream_get_item_for_stream_id(virtual_id, &app_state.config, target, None) { - let input_id: u16 = pli.input_id; - if let Some(input) = app_state.config.get_input_by_id(input_id) { - let mut query_path = if stream_req.action_path.is_empty() { String::new() } else { format!("{}/", stream_req.action_path) }; - query_path = format!("{query_path}{}{stream_ext}", pli.provider_id); - if let Some(stream_url) = get_xtream_player_api_stream_url(input, stream_req.context.to_string().as_str(), query_path.as_str()) { - if user.proxy == ProxyType::Redirect { - debug!("Redirecting stream request to {stream_url}"); - return HttpResponse::Found().insert_header(("Location", stream_url)).finish(); - } - return stream_response(&stream_url, req, Some(input)).await; - } - error!("Cant find stream url for target {target_name}, context {}, stream_id {virtual_id}", stream_req.context); - } else { - error!("Cant find input for target {target_name}, context {}, stream_id {virtual_id}", stream_req.context); - } - } else { - error!("Failed to read xtream item for stream id {}", virtual_id); - } - } else { - debug!("Target has no xtream output {}", target_name); - } - } else { - debug!("Could not find any user {}", stream_req.username); + let (user, target) = try_option_bad_request!(get_user_target_by_credentials(stream_req.username, stream_req.password, api_req, app_state), false, format!("Could not find any user {}", stream_req.username)); + let target_name = &target.name; + if !target.has_output(&TargetType::Xtream) { + debug!("Target has no xtream output {}", target_name); + return HttpResponse::BadRequest().finish(); } - HttpResponse::BadRequest().finish() + let (action_stream_id, stream_ext) = xtream_api_request_separate_number_and_rest(stream_req.stream_id); + let virtual_id: u32 = try_result_bad_request!(action_stream_id.trim().parse()); + let pli = try_result_bad_request!(xtream_repository::xtream_get_item_for_stream_id(virtual_id, &app_state.config, target, None), true, format!("Failed to read xtream item for stream id {}", virtual_id)); + let input = try_option_bad_request!(app_state.config.get_input_by_id(pli.input_id), true, format!("Cant find input for target {target_name}, context {}, stream_id {virtual_id}", stream_req.context)); + + let query_path = if stream_req.action_path.is_empty() { + format!("{}{stream_ext}", pli.provider_id) + } else { + format!("{}/{}{stream_ext}", stream_req.action_path, pli.provider_id) + }; + + let stream_url = try_option_bad_request!(get_xtream_player_api_stream_url(input, stream_req.context.to_string().as_str(), &query_path), true, format!("Cant find stream url for target {target_name}, context {}, stream_id {virtual_id}", stream_req.context)); + + if user.proxy == ProxyType::Redirect { + debug!("Redirecting stream request to {stream_url}"); + return HttpResponse::Found().insert_header(("Location", stream_url)).finish(); + } + + stream_response(&stream_url, req, Some(input)).await } + async fn xtream_player_api_live_stream( req: HttpRequest, api_req: web::Query, @@ -286,8 +314,7 @@ async fn xtream_get_stream_info_content(info_url: &str, input: &ConfigInput) -> async fn xtream_get_stream_info(config: &Config, input: &ConfigInput, target: &ConfigTarget, pli: &XtreamPlaylistItem, info_url: &str, cluster: XtreamCluster) -> Result { if cluster == XtreamCluster::Series { - // TODO if expired then update content ! - if let Ok(content) = xtream_repository::xtream_load_series_info(config, target.name.as_str(), pli.virtual_id) { + if let Some(content) = xtream_repository::xtream_load_series_info(config, target.name.as_str(), pli.virtual_id) { return Ok(content); } } @@ -389,60 +416,33 @@ async fn xtream_player_api_handle_content_action(config: &Config, target_name: & } async fn xtream_get_catchup_response(app_state: &AppState, target: &ConfigTarget, stream_id: &str, start: &str, end: &str) -> HttpResponse { - let virtual_id: u32 = match FromStr::from_str(stream_id) { - Ok(id) => id, - Err(_) => return HttpResponse::BadRequest().finish() - }; + let virtual_id: u32 = try_result_bad_request!(FromStr::from_str(stream_id)); + let pli = try_result_bad_request!(xtream_repository::xtream_get_item_for_stream_id(virtual_id, &app_state.config, target, Some(XtreamCluster::Live))); + let input = try_option_bad_request!(app_state.config.get_input_by_id(pli.input_id)); + let info_url = try_option_bad_request!(get_xtream_player_api_action_url(input, "get_simple_data_table").map(|action_url| format!("{action_url}&stream_id={}&start={start}&end={end}", pli.provider_id))); + let content = try_result_bad_request!(xtream_get_stream_info_content(info_url.as_str(), input).await); + let mut doc: Map = try_result_bad_request!(serde_json::from_str(&content)); + let epg_listings = try_option_bad_request!(doc.get_mut("epg_listings").and_then(Value::as_array_mut)); + let target_path = try_option_bad_request!(get_target_storage_path(&app_state.config, target.name.as_str())); + let mut target_id_mapping = TargetIdMapping::new(&target_path); - !!!! why new ids for catchup ? - - if let Ok(pli) = xtream_repository::xtream_get_item_for_stream_id(virtual_id, &app_state.config, target, Some(XtreamCluster::Live)) { - let input_id = pli.input_id; - if let Some(input) = app_state.config.get_input_by_id(input_id) { - if let Some(info_url) = get_xtream_player_api_action_url(input, "get_simple_data_table") - .map(|action_url| format!("{action_url}&stream_id={}&start={start}&end={end}", pli.provider_id)) { - if let Ok(content) = xtream_get_stream_info_content(info_url.as_str(), input).await { - if let Ok(mut doc) = serde_json::from_str::>(content.as_str()) { - if let Some(epg_listings) = doc.get_mut("epg_listings").and_then(|epg_listings| epg_listings.as_array_mut()) { - match xtream_repository::xtream_load_catchup_id_mapping(&app_state.config, target.name.as_str()) { - Ok(mut mapping) => { - let mut max_id = mapping.max_id(); - for epg_list_value in epg_listings { - if let Some(epg_list_item) = epg_list_value.as_object_mut() { - // TODO epg_id - if let Some(Some(provider_id)) = epg_list_item.get("id").map(|v| v.as_str()) { - if let Ok(provider_stream_id) = u32::from_str(provider_id) { - let stream_id = match mapping.query(provider_stream_id) { - None => { - max_id += 1; - mapping.insert(max_id, provider_stream_id); - max_id - } - Some(mapped_id) => *mapped_id - }; - epg_list_item.insert("id".to_string(), Value::String(stream_id.to_string())); - } - } - } - } - if let Err(err) = mapping.persist() { - error!("Failed to write catchup id mapping {err}"); - return HttpResponse::BadRequest().finish(); - } - } - Err(err) => { error!("Failed to load catchup id mapping {err}"); } - } - } - - if let Ok(result) = serde_json::to_string(&doc) { - return HttpResponse::Ok().content_type(mime::APPLICATION_JSON).body(result); - } - } - } - } + for epg_list_item in epg_listings.iter_mut().filter_map(Value::as_object_mut) { + // TODO epg_id + if let Some(catchup_provider_id) = epg_list_item.get("id").and_then(Value::as_str).and_then(|id| id.parse::().ok()) { + let uuid = hash_string(&format!("{}/{}", pli.url, catchup_provider_id)); + let virtual_id = target_id_mapping.insert_entry(uuid, catchup_provider_id, &PlaylistItemType::Catchup, pli.provider_id); + epg_list_item.insert("id".to_string(), Value::String(virtual_id.to_string())); } } - HttpResponse::BadRequest().finish() + if let Err(err) = target_id_mapping.persist() { + error!("Failed to write catchup id mapping {err}"); + return HttpResponse::BadRequest().finish(); + } + + match serde_json::to_string(&doc) { + Ok(result) => HttpResponse::Ok().content_type(mime::APPLICATION_JSON).body(result), + Err(_) => HttpResponse::BadRequest().finish(), + } } async fn xtream_player_api( diff --git a/src/filter.rs b/src/filter.rs index 2502d57bb..0e6b1141e 100644 --- a/src/filter.rs +++ b/src/filter.rs @@ -196,7 +196,8 @@ impl std::fmt::Display for Filter { PlaylistItemType::Live => "live", PlaylistItemType::Video => "vod", PlaylistItemType::Series => "series", - PlaylistItemType::SeriesInfo => "series" // yes series-info is handled as series in filter + PlaylistItemType::SeriesInfo => "series", // yes series-info is handled as series in filter + _ => "unsupported" }) } Filter::Group(stmt) => { diff --git a/src/model/config.rs b/src/model/config.rs index 5761cc34b..1e9f84b4f 100644 --- a/src/model/config.rs +++ b/src/model/config.rs @@ -24,6 +24,7 @@ use crate::utils::{config_reader, file_utils}; use crate::utils::default_utils::{default_as_default, default_as_false, default_as_true, default_as_empty_list, default_as_frm, default_as_empty_map, default_as_zero_u8, default_as_two_u16}; +use crate::utils::file_lock_manager::FileLockManager; pub(crate) const MAPPER_ATTRIBUTE_FIELDS: &[&str] = &[ "name", "title", "group", "id", "chno", "logo", @@ -809,7 +810,8 @@ pub(crate) struct Config { pub t_sources_file_path: String, #[serde(skip_serializing, skip_deserializing)] pub t_api_proxy_file_path: String, - + #[serde(skip)] + pub file_locks : Arc } impl Config { diff --git a/src/model/playlist.rs b/src/model/playlist.rs index 87a67f096..692841cf4 100644 --- a/src/model/playlist.rs +++ b/src/model/playlist.rs @@ -73,7 +73,7 @@ impl TryFrom for XtreamCluster { PlaylistItemType::Live => Ok(XtreamCluster::Live), PlaylistItemType::Video => Ok(XtreamCluster::Video), PlaylistItemType::Series => Ok(XtreamCluster::Series), - PlaylistItemType::SeriesInfo => Err(format!("Cant convert {item_type}")), + _ => Err(format!("Cant convert {item_type}")), } } } @@ -83,8 +83,10 @@ impl TryFrom for XtreamCluster { pub(crate) enum PlaylistItemType { Live = 1, Video = 2, - Series = 3, - SeriesInfo = 4, + Series = 3, // xtream series description + SeriesInfo = 4, // xtream series info fetched for series description + SeriesEpisode = 5, // from SeriesInfo parsed episodes + Catchup = 6, } impl Default for PlaylistItemType { @@ -110,6 +112,8 @@ impl Display for PlaylistItemType { PlaylistItemType::Video => "video", PlaylistItemType::Series => "series", PlaylistItemType::SeriesInfo => "series-info", + PlaylistItemType::SeriesEpisode => "series-episode", + PlaylistItemType::Catchup => "catchup" }) } } diff --git a/src/processing/playlist_processor.rs b/src/processing/playlist_processor.rs index c6d9b4d2b..7244f7ce1 100644 --- a/src/processing/playlist_processor.rs +++ b/src/processing/playlist_processor.rs @@ -466,7 +466,8 @@ fn flatten_groups(mut playlistgroups: Vec) -> Vec } async fn process_playlist<'a>(playlists: &mut [FetchedPlaylist<'a>], - target: &ConfigTarget, cfg: &Config, + target: &ConfigTarget, + cfg: &Config, stats: &mut HashMap, errors: &mut Vec) -> Result<(), Vec> { let pipe = get_processing_pipe(target); diff --git a/src/repository/bplustree.rs b/src/repository/bplustree.rs index 2dc7a1bf7..5ed74b1a0 100644 --- a/src/repository/bplustree.rs +++ b/src/repository/bplustree.rs @@ -3,9 +3,9 @@ use std::fs::{File, OpenOptions}; use std::io::{self, Read, Seek, SeekFrom, Write}; use std::marker::PhantomData; use std::path::Path; + use flate2::Compression; use log::error; - use serde::{Deserialize, Serialize}; const BINCODE_OVERHEAD: usize = 4; @@ -17,6 +17,18 @@ fn is_multiple_of_block_size(file: &File) -> io::Result { Ok(file_size % (BLOCK_SIZE as u64) == 0) // Check if file size is a multiple of BLOCK_SIZE } +fn is_file_valid(file: File) -> io::Result { + match is_multiple_of_block_size(&file) { + Ok(valid) => { + if !valid { + return Err(io::Error::new(io::ErrorKind::InvalidData, format!("Tree file has to be multiple of block size {BLOCK_SIZE}"))); + } + } + Err(err) => return Err(err) + } + Ok(file) +} + #[inline] fn u32_from_bytes(bytes: &[u8]) -> io::Result { Ok(u32::from_le_bytes(bytes.try_into().map_err(|e: TryFromSliceError| io::Error::new(io::ErrorKind::Other, e.to_string()))?)) @@ -39,6 +51,25 @@ where } +fn get_entry_index_upper_bound(keys: &Vec, key: &K) -> usize +where + K: Ord + Serialize + for<'de> Deserialize<'de> + Clone, + V: Serialize + for<'de> Deserialize<'de> + Clone, +{ + let mut left = 0; + let mut right = keys.len(); + while left < right { + let mid = left + ((right - left) >> 1); + if &keys[mid] <= key { + left = mid + 1; + } else { + right = mid; + } + } + left +} + + #[derive(Serialize, Deserialize, Debug, Clone)] struct BPlusTreeNode { keys: Vec, @@ -81,6 +112,7 @@ where } } + #[allow(dead_code)] fn query(&self, key: &K) -> Option<&V> { if self.is_leaf { return match self.keys.binary_search(&key) { @@ -113,17 +145,7 @@ where } fn get_entry_index_upper_bound(&self, key: &K) -> usize { - let mut left = 0; - let mut right = self.keys.len(); - while left < right { - let mid = left + ((right - left) >> 1); - if &self.keys[mid] <= key { - left = mid + 1; - } else { - right = mid; - } - } - left + get_entry_index_upper_bound::(&self.keys, key) } fn insert(&mut self, key: K, v: V, inner_order: usize, leaf_order: usize) -> Option> { @@ -181,7 +203,7 @@ where self.children.iter().for_each(|child| child.traverse(visit)); } - fn serialize_to_blocks(&self, file: &mut W, buffer: &mut Vec, offset: u64) -> io::Result { + fn serialize_to_block(&self, file: &mut W, buffer: &mut Vec, offset: u64) -> io::Result { let mut current_offset = offset; let buffer_slice = &mut buffer[..]; @@ -204,7 +226,7 @@ where encoder.write_all(&values_encoded)?; let compressed_bytes = encoder.finish()?; let values_bytes_len = compressed_bytes.len(); - buffer_slice[write_pos..write_pos + 4].copy_from_slice(&(values_bytes_len as u32).to_le_bytes()); + buffer_slice[write_pos..write_pos + 4].copy_from_slice(&(values_bytes_len as u32).to_le_bytes()); write_pos += 4; buffer_slice[write_pos..write_pos + values_bytes_len].copy_from_slice(&compressed_bytes); write_pos += values_bytes_len; @@ -220,7 +242,7 @@ where let mut pointer = Vec::with_capacity(self.children.len()); for child in &self.children { pointer.push(current_offset); - current_offset = child.serialize_to_blocks(file, buffer, current_offset)?; + current_offset = child.serialize_to_block(file, buffer, current_offset)?; } let pointer_encoded = bincode_serialize(&pointer)?; @@ -234,7 +256,7 @@ where Ok(current_offset) } - fn deserialize_from_blocks(file: &mut R, buffer: &mut Vec, offset: u64, nested: bool) -> io::Result<(Self, Option>)> { + fn deserialize_from_block(file: &mut R, buffer: &mut Vec, offset: u64, nested: bool) -> io::Result<(Self, Option>)> { file.seek(SeekFrom::Start(offset))?; file.read_exact(buffer)?; @@ -272,7 +294,7 @@ where let nodes: Result>, io::Error> = pointers .iter() .map(|pointer| { - BPlusTreeNode::::deserialize_from_blocks(file, buffer, *pointer, nested) + BPlusTreeNode::::deserialize_from_block(file, buffer, *pointer, nested) .map(|(node, _)| node) .map_err(|err| io::Error::new(io::ErrorKind::Other, err.to_string())) }) @@ -352,6 +374,7 @@ where } } + #[allow(dead_code)] pub(crate) fn query(&self, key: &K) -> Option<&V> { self.root.query(key) } @@ -359,7 +382,7 @@ where pub(crate) fn serialize(&self, filepath: &Path) -> io::Result { let mut file = OpenOptions::new().write(true).create(true).open(filepath)?; let mut buffer = vec![0u8; BLOCK_SIZE]; - let result = self.root.serialize_to_blocks(&mut file, &mut buffer, 0u64); + let result = self.root.serialize_to_block(&mut file, &mut buffer, 0u64); file.flush()?; result } @@ -370,7 +393,7 @@ where return Err(std::io::Error::new(std::io::ErrorKind::InvalidData, format!("Tree file has to be multiple of block size {BLOCK_SIZE}"))); } let mut buffer = vec![0u8; BLOCK_SIZE]; - let (root, _) = BPlusTreeNode::deserialize_from_blocks(&mut file, &mut buffer, 0, true)?; + let (root, _) = BPlusTreeNode::deserialize_from_block(&mut file, &mut buffer, 0, true)?; Ok(BPlusTree::new_with_root(root)) } @@ -382,6 +405,33 @@ where } } +fn query_tree(file: &mut File, key: &K) -> Option +where + K: Ord + Serialize + for<'de> Deserialize<'de> + Clone, + V: Serialize + for<'de> Deserialize<'de> + Clone, +{ + let mut offset = 0; + let mut buffer = vec![0u8; BLOCK_SIZE]; + loop { + match BPlusTreeNode::::deserialize_from_block(file, &mut buffer, offset, false) { + Ok((node, pointers)) => { + if node.is_leaf { + return match node.keys.binary_search(key) { + Ok(idx) => node.values.get(idx).cloned(), + Err(_) => None, + }; + } + let child_idx = get_entry_index_upper_bound::(&node.keys, key); + offset = *pointers.unwrap().get(child_idx).unwrap(); + } + Err(err) => { + error!("Failed to read id tree from file {err}"); + return None; + } + }; + } +} + pub(crate) struct BPlusTreeQuery { file: File, _marker_k: PhantomData, @@ -394,15 +444,7 @@ where V: Serialize + for<'de> Deserialize<'de> + Clone, { pub(crate) fn try_new(filepath: &Path) -> io::Result { - let file = File::open(filepath)?; - match is_multiple_of_block_size(&file) { - Ok(valid) => { - if !valid { - return Err(std::io::Error::new(std::io::ErrorKind::InvalidData, format!("Tree file has to be multiple of block size {BLOCK_SIZE}"))); - } - } - Err(err) => return Err(err) - } + let file = is_file_valid(File::open(filepath)?)?; Ok(BPlusTreeQuery { file, _marker_k: Default::default(), @@ -410,44 +452,81 @@ where }) } - fn get_entry_index_upper_bound(keys: &Vec, key: &K) -> usize { - let mut left = 0; - let mut right = keys.len(); - while left < right { - let mid = left + ((right - left) >> 1); - if &keys[mid] <= key { - left = mid + 1; - } else { - right = mid; - } + pub(crate) fn query(&mut self, key: &K) -> Option { + query_tree(&mut self.file, key) + } +} + +pub(crate) struct BPlusTreeUpdate { + file: File, + _marker_k: PhantomData, + _marker_v: PhantomData, +} + +impl BPlusTreeUpdate +where + K: Ord + Serialize + for<'de> Deserialize<'de> + Clone, + V: Serialize + for<'de> Deserialize<'de> + Clone, +{ + pub(crate) fn try_new(filepath: &Path) -> io::Result { + if !filepath.exists() { + return Err(io::Error::new(io::ErrorKind::NotFound, format!("File not found {}", filepath.to_str().unwrap_or("?")))); } - left + let file = OpenOptions::new() + .write(true) + .read(true) + .open(filepath)?; + let file = is_file_valid(file)?; + Ok(BPlusTreeUpdate { + file, + _marker_k: Default::default(), + _marker_v: Default::default(), + }) } pub(crate) fn query(&mut self, key: &K) -> Option { + query_tree(&mut self.file, key) + } + + fn serialize_node(&mut self, offset: u64, node: BPlusTreeNode) -> io::Result { + let mut buffer = vec![0u8; BLOCK_SIZE]; + let result = node.serialize_to_block(&mut self.file, &mut buffer, offset); + self.file.flush()?; + result + } + + pub(crate) fn update(&mut self, key: &K, value: V) -> io::Result { let mut offset = 0; let mut buffer = vec![0u8; BLOCK_SIZE]; loop { - match BPlusTreeNode::::deserialize_from_blocks(&mut self.file, &mut buffer, offset, false) { - Ok((node, pointers)) => { + match BPlusTreeNode::::deserialize_from_block(&mut self.file, &mut buffer, offset, false) { + Ok((mut node, pointers)) => { if node.is_leaf { return match node.keys.binary_search(key) { - Ok(idx) => node.values.get(idx).cloned(), - Err(_) => None, + Ok(idx) => { + let old_value = node.values.get(idx); + if let Some(_) = old_value { + node.values[idx] = value; + return self.serialize_node(offset, node); + } + Err(io::Error::new(io::ErrorKind::NotFound, "Entry not found")) + } + Err(_) => Err(io::Error::new(io::ErrorKind::NotFound, "Entry not found")), }; } - let child_idx = BPlusTreeQuery::::get_entry_index_upper_bound(&node.keys, key); + let child_idx = get_entry_index_upper_bound::(&node.keys, key); offset = *pointers.unwrap().get(child_idx).unwrap(); } Err(err) => { error!("Failed to read id tree from file {err}"); - return None; + return Err(io::Error::new(io::ErrorKind::NotFound, format!("Failed to read id tree from file {err}"))); } }; } } } + #[cfg(test)] mod tests { use std::io; @@ -455,7 +534,7 @@ mod tests { use serde::{Deserialize, Serialize}; - use crate::repository::bplustree::{BPlusTree, BPlusTreeQuery}; + use crate::repository::bplustree::{BPlusTree, BPlusTreeQuery, BPlusTreeUpdate}; // Example usage with a simple struct #[derive(Serialize, Deserialize, Debug, Clone, PartialEq, Eq)] @@ -507,6 +586,31 @@ mod tests { }), "Entry {} not found", i); } + let mut tree_update: BPlusTreeUpdate = BPlusTreeUpdate::try_new(&filepath)?; + for i in 0u32..=500 { + if let Some(record) = tree_update.query(&i) { + let new_record = Record { + id: record.id, + data: format!("Entry {}", record.id + 9000), + }; + tree_update.update(&i, new_record)?; + } else { + assert!(false, "Entry {} not found", i); + } + } + + let mut tree_query: BPlusTreeQuery = BPlusTreeQuery::try_new(&filepath)?; + for i in 0u32..=500 { + let found = tree_query.query(&i); + assert!(found.is_some(), "Entry {} not found", i); + let entry = found.unwrap(); + let expected = Record { + id: i, + data: format!("Entry {}", i + 9000), + }; + assert!(entry.eq(&expected), "Entry not equal {:?} != {:?}", entry, expected); + } + Ok(()) } } diff --git a/src/repository/id_mapping.rs b/src/repository/id_mapping.rs index 7ad77478a..a64199b0a 100644 --- a/src/repository/id_mapping.rs +++ b/src/repository/id_mapping.rs @@ -1,3 +1,5 @@ +// This module is not included in the build + use std::cmp::max; use std::io::Error; use std::path::{Path, PathBuf}; diff --git a/src/repository/m3u_repository.rs b/src/repository/m3u_repository.rs index 5e1467612..b6028de9d 100644 --- a/src/repository/m3u_repository.rs +++ b/src/repository/m3u_repository.rs @@ -4,15 +4,15 @@ use std::path::{Path, PathBuf}; use log::error; -use crate::{create_m3u_filter_error}; use crate::api::api_utils::get_user_server_info; +use crate::create_m3u_filter_error; use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; use crate::model::api_proxy::{ProxyType, ProxyUserCredentials}; use crate::model::config::{Config, ConfigTarget}; use crate::model::playlist::{M3uPlaylistItem, PlaylistGroup, PlaylistItem, PlaylistItemType}; -use crate::repository::indexed_document_reader::{IndexedDocumentReader}; +use crate::repository::indexed_document_reader::IndexedDocumentReader; use crate::repository::indexed_document_writer::IndexedDocumentWriter; -use crate::repository::storage::{ensure_target_storage_path}; +use crate::repository::storage::ensure_target_storage_path; use crate::utils::file_utils; macro_rules! cant_write_result { @@ -64,25 +64,26 @@ fn persist_m3u_playlist_as_text(target: &ConfigTarget, cfg: &Config, m3u_playlis pub(crate) fn m3u_write_playlist(target: &ConfigTarget, cfg: &Config, target_path: &Path, new_playlist: &[PlaylistGroup]) -> Result<(), M3uFilterError> { if !new_playlist.is_empty() { if let Some((m3u_path, idx_path)) = m3u_get_file_paths(target_path) { - let m3u_playlist = new_playlist.iter() .flat_map(|pg| &pg.channels) .filter(|&pli| pli.header.borrow().item_type != PlaylistItemType::SeriesInfo) .map(PlaylistItem::to_m3u).collect::>(); persist_m3u_playlist_as_text(target, cfg, &m3u_playlist); - - match IndexedDocumentWriter::new(m3u_path.clone(), idx_path) { - Ok(mut writer) => { - for m3u in m3u_playlist { - match writer.write_doc(m3u.virtual_id, &m3u) { - Ok(_) => {}, - Err(err) => return Err(cant_write_result!(&m3u_path, err)) + { + let _file_lock = cfg.file_locks.write_lock(&m3u_path).map_err(|err| M3uFilterError::new(M3uFilterErrorKind::Info, format!("{}", err)))?; + match IndexedDocumentWriter::new(m3u_path.clone(), idx_path) { + Ok(mut writer) => { + for m3u in m3u_playlist { + match writer.write_doc(m3u.virtual_id, &m3u) { + Ok(_) => {} + Err(err) => return Err(cant_write_result!(&m3u_path, err)) + } } + writer.flush().map_err(|err| cant_write_result!(&m3u_path, err))?; } - writer.flush().map_err(|err| cant_write_result!(&m3u_path, err))?; + Err(err) => return Err(cant_write_result!(&m3u_path, err)) } - Err(err) => return Err(cant_write_result!(&m3u_path, err)) } } } @@ -93,36 +94,43 @@ pub(crate) fn m3u_load_rewrite_playlist(cfg: &Config, target: &ConfigTarget, use match ensure_target_storage_path(cfg, target.name.as_str()) { Ok(target_path) => { if let Some((m3u_path, _)) = m3u_get_file_paths(&target_path) { - match IndexedDocumentReader::::new(&m3u_path) { - Ok(mut reader) => { - let server_info = get_user_server_info(cfg, user); - let url = format!("{}/m3u-stream/{}/{}", server_info.get_base_url(), user.username, user.password); - let mut result = vec![]; - result.push("#EXTM3U".to_string()); - for m3u_pli in reader.by_ref() { - match user.proxy { - ProxyType::Reverse => { - result.push(m3u_pli.to_m3u(target, Some(format!("{url}/{}", m3u_pli.virtual_id).as_str()))); - } - ProxyType::Redirect => { - result.push(m3u_pli.to_m3u(target, None)); + { + let _file_lock = cfg.file_locks.read_lock(&m3u_path).map_err(|err| { + error!("Could not lock document {:?}: {}", m3u_path, err); + Error::new(ErrorKind::Other, format!("Document Reader error for target {}", &target.name)) + }).ok()?; + + match IndexedDocumentReader::::new(&m3u_path) { + Ok(mut reader) => { + let server_info = get_user_server_info(cfg, user); + let url = format!("{}/m3u-stream/{}/{}", server_info.get_base_url(), user.username, user.password); + let mut result = vec![]; + result.push("#EXTM3U".to_string()); + for m3u_pli in reader.by_ref() { + match user.proxy { + ProxyType::Reverse => { + result.push(m3u_pli.to_m3u(target, Some(format!("{url}/{}", m3u_pli.virtual_id).as_str()))); + } + ProxyType::Redirect => { + result.push(m3u_pli.to_m3u(target, None)); + } } + }; + if reader.by_ref().has_error() { + error!("Could not deserialize m3u item {}", &m3u_path.to_str().unwrap()); + } else { + return Some(result.join("\n")); } - }; - if reader.by_ref().has_error() { - error!("Could not deserialize m3u item {}", &m3u_path.to_str().unwrap()); - } else { - return Some(result.join("\n")); } - } - Err(err) => { - error!("Could not deserialize file {} - {}", &m3u_path.to_str().unwrap(), err); + Err(err) => { + error!("Could not deserialize file {} - {}", &m3u_path.to_str().unwrap(), err); + } } } } else { error!("Could not open files for target {}", &target.name); } - }, + } Err(err) => { error!("Could not find storage path for target {} - {}", target.name.as_str(), err); } @@ -130,9 +138,12 @@ pub(crate) fn m3u_load_rewrite_playlist(cfg: &Config, target: &ConfigTarget, use None } -pub(crate) fn m3u_get_item_for_stream_id(stream_id: u32, m3u_path: &Path, idx_path: &Path) -> Result { +pub(crate) fn m3u_get_item_for_stream_id(cfg: &Config, stream_id: u32, m3u_path: &Path, idx_path: &Path) -> Result { if stream_id < 1 { return Err(Error::new(ErrorKind::Other, "id should start with 1")); } - IndexedDocumentReader::::read_indexed_item(m3u_path, idx_path, stream_id) + { + let _file_lock = cfg.file_locks.read_lock(&m3u_path)?; + IndexedDocumentReader::::read_indexed_item(m3u_path, idx_path, stream_id) + } } \ No newline at end of file diff --git a/src/repository/mod.rs b/src/repository/mod.rs index 91382caaf..7db3f6877 100644 --- a/src/repository/mod.rs +++ b/src/repository/mod.rs @@ -8,5 +8,4 @@ pub(crate) mod storage; mod indexed_document_writer; mod indexed_document_reader; pub(crate) mod target_id_mapping; -pub(crate) mod id_mapping; pub(crate) mod bplustree; diff --git a/src/repository/playlist_repository.rs b/src/repository/playlist_repository.rs index 38767d3b2..298369eed 100644 --- a/src/repository/playlist_repository.rs +++ b/src/repository/playlist_repository.rs @@ -15,54 +15,61 @@ pub(crate) fn persist_playlist(playlist: &mut [PlaylistGroup], epg: Option<&Epg> match ensure_target_storage_path(cfg, target.name.as_str()) { Ok(target_path) => { - // reassign old virtual id or assign a new one - let mut target_id_mapping = TargetIdMapping::new(&get_target_id_mapping_file(&target_path)); - for group in &mut *playlist { - for channel in &group.channels { - let mut header = channel.header.borrow_mut(); - match header.get_provider_id() { - Some(provider_id) => { - let uuid = header.get_uuid(); - let item_type = header.item_type; - header.virtual_id = match target_id_mapping.get_by_uuid(&header.uuid) { - None => target_id_mapping.insert_entry(provider_id, **uuid, &item_type, 0), - Some(existing_id) => *existing_id - }; + { + let target_id_mapping_file = get_target_id_mapping_file(&target_path); + match cfg.file_locks.write_lock(&target_id_mapping_file) { + Ok(_file_lock) => { + // reassign old virtual id or assign a new one + let mut target_id_mapping = TargetIdMapping::new(&target_id_mapping_file); + for group in &mut *playlist { + for channel in &group.channels { + let mut header = channel.header.borrow_mut(); + match header.get_provider_id() { + Some(provider_id) => { + let uuid = header.get_uuid(); + let item_type = header.item_type; + header.virtual_id = target_id_mapping.insert_entry(**uuid, provider_id, &item_type, 0); + } + None => { + errors.push(M3uFilterError::new(M3uFilterErrorKind::Info, format!("Playlistitem has no provider id: {}", &header.title))); + } + } + } } - None => { - errors.push(M3uFilterError::new(M3uFilterErrorKind::Info, format!("Playlistitem has no provider id: {}", &header.title))); - } - } - } - } - for output in &target.output { - match match output.target { - TargetType::M3u => m3u_write_playlist(target, cfg, &target_path, playlist), - TargetType::Xtream => xtream_write_playlist(target, cfg, playlist), - TargetType::Strm => kodi_write_strm_playlist(target, cfg, playlist, &output.filename), - } { - Ok(()) => { - if let Err(err) = target_id_mapping.persist() { - errors.push(M3uFilterError::new(M3uFilterErrorKind::Info, err.to_string())); - } - if !playlist.is_empty() { - match epg_write(target, cfg, &target_path, epg, output) { - Ok(()) => {} + for output in &target.output { + match match output.target { + TargetType::M3u => m3u_write_playlist(target, cfg, &target_path, playlist), + TargetType::Xtream => xtream_write_playlist(target, cfg, playlist), + TargetType::Strm => kodi_write_strm_playlist(target, cfg, playlist, &output.filename), + } { + Ok(()) => { + if let Err(err) = target_id_mapping.persist() { + errors.push(M3uFilterError::new(M3uFilterErrorKind::Info, err.to_string())); + } + if !playlist.is_empty() { + match epg_write(target, cfg, &target_path, epg, output) { + Ok(()) => {} + Err(err) => errors.push(err) + } + } + } Err(err) => errors.push(err) } } + if let Err(err) = target_id_mapping.persist() { + errors.push(M3uFilterError::new(M3uFilterErrorKind::Info, err.to_string())); + } + } + Err(err) => { + errors.push(M3uFilterError::new(M3uFilterErrorKind::Info, err.to_string())); } - Err(err) => errors.push(err) } } - if let Err(err) = target_id_mapping.persist() { - errors.push(M3uFilterError::new(M3uFilterErrorKind::Info, err.to_string())); - } } Err(err) => { errors.push(err); - }, + } }; if errors.is_empty() { Ok(()) } else { Err(errors) } diff --git a/src/repository/storage.rs b/src/repository/storage.rs index 62e899703..a2ed476b9 100644 --- a/src/repository/storage.rs +++ b/src/repository/storage.rs @@ -8,7 +8,7 @@ pub(crate) fn hash_string(url: &str) -> [u8; 32] { hash.into() // Konvertiere den Hash in ein Array mit fester Größe } -pub(crate) fn get_target_id_mapping_file(target_path: &Path) -> PathBuf { +pub(in crate::repository) fn get_target_id_mapping_file(target_path: &Path) -> PathBuf { target_path.join(PathBuf::from("id_mapping.db")) } diff --git a/src/repository/target_id_mapping.rs b/src/repository/target_id_mapping.rs index 2cbdecbbe..5e204c124 100644 --- a/src/repository/target_id_mapping.rs +++ b/src/repository/target_id_mapping.rs @@ -31,6 +31,10 @@ impl VirtualIdRecord { pub(crate) fn is_expired(&self) -> bool { (Local::now().timestamp() - self.last_updated) > EXPIRATION_DURATION } + + pub(crate) fn copy_update_timestamp(&self) -> Self { + Self::new(self.provider_id, self.virtual_id, self.item_type, self.parent_virtual_id, self.uuid) + } } pub(crate) struct TargetIdMapping { @@ -69,16 +73,17 @@ impl TargetIdMapping { } } - pub(crate) fn insert_entry(&mut self, provider_id: u32, uuid: [u8; 32], item_type: &PlaylistItemType, parent_virtual_id: u32) -> u32 { - self.dirty = true; - self.virtual_id_counter += 1; - let record = VirtualIdRecord::new(provider_id, self.virtual_id_counter, *item_type, parent_virtual_id, uuid); - self.by_virtual_id.insert(self.virtual_id_counter, record); - self.virtual_id_counter - } - - pub(crate) fn get_by_uuid(&self, uuid: &[u8; 32]) -> Option<&u32> { - self.by_uuid.get(uuid) + pub(crate) fn insert_entry(&mut self, uuid: [u8; 32], provider_id: u32, item_type: &PlaylistItemType, parent_virtual_id: u32) -> u32 { + match self.by_uuid.get(&uuid) { + None => { + self.dirty = true; + self.virtual_id_counter += 1; + let record = VirtualIdRecord::new(provider_id, self.virtual_id_counter, *item_type, parent_virtual_id, uuid); + self.by_virtual_id.insert(self.virtual_id_counter, record); + self.virtual_id_counter + } + Some(record) => *record + } } pub(crate) fn persist(&mut self) -> Result<(), Error> { diff --git a/src/repository/xtream_repository.rs b/src/repository/xtream_repository.rs index f03b97e34..bd0327aed 100644 --- a/src/repository/xtream_repository.rs +++ b/src/repository/xtream_repository.rs @@ -11,9 +11,8 @@ use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; use crate::model::config::{Config, ConfigTarget}; use crate::model::playlist::{PlaylistGroup, PlaylistItem, PlaylistItemType, XtreamCluster, XtreamPlaylistItem}; use crate::model::xtream::XtreamMappingOptions; -use crate::repository::bplustree::{BPlusTreeQuery}; -use crate::repository::id_mapping::IdMapping; -use crate::repository::indexed_document_reader::{IndexedDocumentReader}; +use crate::repository::bplustree::{BPlusTreeQuery, BPlusTreeUpdate}; +use crate::repository::indexed_document_reader::IndexedDocumentReader; use crate::repository::indexed_document_writer::IndexedDocumentWriter; use crate::repository::storage::{get_target_id_mapping_file, get_target_storage_path, hash_string}; use crate::repository::target_id_mapping::{TargetIdMapping, VirtualIdRecord}; @@ -56,29 +55,28 @@ fn xtream_get_info_file_paths(storage_path: &Path, cluster: XtreamCluster) -> Op None } -fn xtream_get_catchup_id_mapping_file_path(storage_path: &Path) -> PathBuf { - storage_path.join("catchup_mapping.db") -} - -fn write_playlists_to_file(storage_path: &Path, collections: Vec<(XtreamCluster, &mut [PlaylistItem])>) -> Result<(), M3uFilterError> { +fn write_playlists_to_file(cfg: &Config, storage_path: &Path, collections: Vec<(XtreamCluster, &mut [PlaylistItem])>) -> Result<(), M3uFilterError> { for (cluster, playlist) in collections { let (xtream_path, idx_path) = xtream_get_file_paths(storage_path, &cluster); - match IndexedDocumentWriter::new(xtream_path.clone(), idx_path) { - Ok(mut writer) => { - for item in playlist { - match item.to_xtream() { - Ok(xtream) => { - match writer.write_doc(item.header.borrow().virtual_id, &xtream) { - Ok(_) => {} - Err(err) => return Err(cant_write_result!(&xtream_path, err)) + { + let _file_lock = cfg.file_locks.write_lock(&xtream_path).map_err(|err| M3uFilterError::new(M3uFilterErrorKind::Info, format!("{}", err)))?; + match IndexedDocumentWriter::new(xtream_path.clone(), idx_path) { + Ok(mut writer) => { + for item in playlist { + match item.to_xtream() { + Ok(xtream) => { + match writer.write_doc(item.header.borrow().virtual_id, &xtream) { + Ok(_) => {} + Err(err) => return Err(cant_write_result!(&xtream_path, err)) + } } - }, - Err(err) => return Err(cant_write_result!(&xtream_path, err)) + Err(err) => return Err(cant_write_result!(&xtream_path, err)) + } } + writer.flush().map_err(|err| cant_write_result!(&xtream_path, err))?; } - writer.flush().map_err(|err| cant_write_result!(&xtream_path, err))?; + Err(err) => return Err(cant_write_result!(&xtream_path, err)) } - Err(err) => return Err(cant_write_result!(&xtream_path, err)) } } Ok(()) @@ -221,7 +219,7 @@ pub(crate) fn xtream_write_playlist(target: &ConfigTarget, cfg: &Config, playlis } } - match write_playlists_to_file(&path, vec![ + match write_playlists_to_file(cfg, &path, vec![ (XtreamCluster::Live, &mut live_col), (XtreamCluster::Video, &mut vod_col), (XtreamCluster::Series, &mut series_col)]) { @@ -250,186 +248,222 @@ pub(crate) fn xtream_get_collection_path(cfg: &Config, target_name: &str, collec Err(Error::new(ErrorKind::Other, format!("Cant find collection: {target_name}/{collection_name}"))) } -fn xtream_read_item_for_stream_id(stream_id: u32, storage_path: &Path, cluster: &XtreamCluster) -> Result { +fn xtream_read_item_for_stream_id(cfg: &Config, stream_id: u32, storage_path: &Path, cluster: &XtreamCluster) -> Result { let (xtream_path, idx_path) = xtream_get_file_paths(storage_path, cluster); - return IndexedDocumentReader::::read_indexed_item(&xtream_path, &idx_path, stream_id); -} - -fn xtream_read_series_item_for_stream_id(stream_id: u32, storage_path: &Path) -> Result { - let (xtream_path, idx_path) = xtream_get_file_paths_for_series(storage_path); - return IndexedDocumentReader::::read_indexed_item(&xtream_path, &idx_path, stream_id); -} - -pub(crate) fn xtream_get_item_for_stream_id(virtual_id: u32, config: &Config, target: &ConfigTarget, xtream_cluster: Option) -> Result { - if let Some(target_path) = get_target_storage_path(config, target.name.as_str()) { - if let Some(storage_path) = xtream_get_storage_path(config, target.name.as_str()) { - match BPlusTreeQuery::::try_new(&get_target_id_mapping_file(&target_path)) { - Ok(mut target_id_mapping) => { - return match target_id_mapping.query(&virtual_id) { - Some(mapping) => { - if mapping.item_type == PlaylistItemType::SeriesInfo { - xtream_read_series_item_for_stream_id(virtual_id, &storage_path) - } else if mapping.item_type == PlaylistItemType::Series && mapping.parent_virtual_id > 0 { - // we load the original series item - match xtream_read_series_item_for_stream_id(mapping.parent_virtual_id, &storage_path) { - Ok(mut item) => { - // we need to replace the provider id with the episode provider id - item.provider_id = mapping.provider_id; - Ok(item) - }, - Err(err) => Err(err) - } - } else { - let cluster = match xtream_cluster { - Some(c) => Some(c), - None => match XtreamCluster::try_from(mapping.item_type) { - Ok(item_type) => Some(item_type), - Err(_) => None - } - }; - match cluster { - Some(xc) => xtream_read_item_for_stream_id(virtual_id, &storage_path, &xc), - None => Err(Error::new(ErrorKind::Other, format!("Could not determine cluster for xtream item with stream-id {virtual_id}"))) - } - } - }, - None => Err(Error::new(ErrorKind::Other, format!("Could not find mappping for target {} and id {}", target.name, virtual_id))), - }; - }, - Err(err) => return Err(Error::new(ErrorKind::Other, format!("Could not load id mappping for target {} err:{}", target.name, err.to_string()))) - }; - } else { - return Err(Error::new(ErrorKind::Other, format!("Could not find path for target {} xtream output", &target.name))); - } - } else { - return Err(Error::new(ErrorKind::Other, format!("Could not find path for target {}", &target.name))); + { + let _file_lock = cfg.file_locks.read_lock(&xtream_path)?; + IndexedDocumentReader::::read_indexed_item(&xtream_path, &idx_path, stream_id) } } +fn xtream_read_series_item_for_stream_id(cfg: &Config, stream_id: u32, storage_path: &Path) -> Result { + let (xtream_path, idx_path) = xtream_get_file_paths_for_series(storage_path); + { + let _file_lock = cfg.file_locks.read_lock(&xtream_path)?; + return IndexedDocumentReader::::read_indexed_item(&xtream_path, &idx_path, stream_id); + } +} + +macro_rules! try_cluster { + ($xtream_cluster:expr, $item_type:expr, $virtual_id:expr) => { + $xtream_cluster.or_else(|| XtreamCluster::try_from($item_type).ok()) + .ok_or_else(|| Error::new(ErrorKind::Other, format!("Could not determine cluster for xtream item with stream-id {}", $virtual_id))) + }; +} + +pub(crate) fn xtream_get_item_for_stream_id( + virtual_id: u32, + config: &Config, + target: &ConfigTarget, + xtream_cluster: Option, +) -> Result { + let target_path = get_target_storage_path(config, target.name.as_str()) + .ok_or_else(|| Error::new(ErrorKind::Other, format!("Could not find path for target {}", &target.name)))?; + let storage_path = xtream_get_storage_path(config, target.name.as_str()) + .ok_or_else(|| Error::new(ErrorKind::Other, format!("Could not find path for target {} xtream output", &target.name)))?; + { + let target_id_mapping_file = get_target_id_mapping_file(&target_path); + let _file_lock = config.file_locks.read_lock(&target_id_mapping_file) + .map_err(|err| Error::new(ErrorKind::Other, format!("Could not get lock for id mapping for target {} err:{}", target.name, err.to_string())))?; + + let mut target_id_mapping = BPlusTreeQuery::::try_new(&target_id_mapping_file) + .map_err(|err| Error::new(ErrorKind::Other, format!("Could not load id mapping for target {} err:{}", target.name, err.to_string())))?; + + let mapping = target_id_mapping + .query(&virtual_id) + .ok_or_else(|| Error::new(ErrorKind::Other, format!("Could not find mapping for target {} and id {}", target.name, virtual_id)))?; + + match mapping.item_type { + PlaylistItemType::SeriesInfo => xtream_read_series_item_for_stream_id(config, virtual_id, &storage_path), + PlaylistItemType::SeriesEpisode => { + let mut item = xtream_read_series_item_for_stream_id(config, mapping.parent_virtual_id, &storage_path)?; + item.provider_id = mapping.provider_id; + Ok(item) + } + PlaylistItemType::Catchup => { + let cluster = try_cluster!(xtream_cluster, mapping.item_type, virtual_id)?; + let mut item = xtream_read_item_for_stream_id(config, mapping.parent_virtual_id, &storage_path, &cluster)?; + item.provider_id = mapping.provider_id; + Ok(item) + } + _ => { + let cluster = try_cluster!(xtream_cluster, mapping.item_type, virtual_id)?; + xtream_read_item_for_stream_id(config, virtual_id, &storage_path, &cluster) + } + } + } +} + + pub(crate) fn xtream_load_rewrite_playlist(cluster: &XtreamCluster, config: &Config, target: &ConfigTarget, category_id: u32) -> Result { if let Some(storage_path) = xtream_get_storage_path(config, target.name.as_str()) { let (xtream_path, _) = xtream_get_file_paths(&storage_path, cluster); - match IndexedDocumentReader::::new(&xtream_path) { - Ok(mut reader) => { - let options = XtreamMappingOptions::from_target_options(target.options.as_ref()); - let result: Vec = reader.by_ref().filter(|pli| category_id == 0 || pli.category_id == category_id) - .map(|pli| pli.to_doc(&options)).collect(); - if reader.by_ref().has_error() { - error!("Could not deserialize item {}", &xtream_path.to_str().unwrap()); - } else { - return Ok(serde_json::to_string(&result).unwrap()); + { + let _file_lock = config.file_locks.read_lock(&xtream_path)?; + match IndexedDocumentReader::::new(&xtream_path) { + Ok(mut reader) => { + let options = XtreamMappingOptions::from_target_options(target.options.as_ref()); + let result: Vec = reader.by_ref().filter(|pli| category_id == 0 || pli.category_id == category_id) + .map(|pli| pli.to_doc(&options)).collect(); + if reader.by_ref().has_error() { + error!("Could not deserialize item {}", &xtream_path.to_str().unwrap()); + } else { + return Ok(serde_json::to_string(&result).unwrap()); + } + } + Err(err) => { + error!("Could not deserialize file {} - {}", &xtream_path.to_str().unwrap(), err); } - } - Err(err) => { - error!("Could not deserialize file {} - {}", &xtream_path.to_str().unwrap(), err); } } } Err(Error::new(ErrorKind::Other, format!("Failed to find xtream storage for target {}", &target.name))) } +macro_rules! try_option_ok { + ($option:expr) => { + match $option { + Some(value) => value, + None => return Ok(()), + } + }; +} pub(crate) fn xtream_write_series_info(config: &Config, target_name: &str, - series_id: u32, + series_info_id: u32, content: &str) -> Result<(), Error> { - if let Some(storage_path) = xtream_get_storage_path(config, target_name) { - if let Some((info_path, idx_path)) = xtream_get_info_file_paths(&storage_path, XtreamCluster::Series) { - return match IndexedDocumentWriter::new_append(info_path.clone(), idx_path) { - Ok(mut writer) => { - match writer.write_doc(series_id, content) { - Ok(_) => {}, - Err(_) => return Err(Error::new(ErrorKind::Other, format!("failed to write xtream series info for target {target_name}"))) - } - return Ok(writer.flush()?); - } - Err(err) => Err(err) - }; - } + let target_path = try_option_ok!(get_target_storage_path(config, target_name)); + let storage_path = try_option_ok!(xtream_get_storage_path(config, target_name)); + let (info_path, idx_path) = try_option_ok!(xtream_get_info_file_paths(&storage_path, XtreamCluster::Series)); + + { + let _file_lock = config.file_locks.write_lock(&info_path)?; + let mut writer = IndexedDocumentWriter::new_append(info_path.clone(), idx_path)?; + + writer + .write_doc(series_info_id, content) + .map_err(|_| Error::new(ErrorKind::Other, format!("failed to write xtream series info for target {target_name}")))?; + + writer.flush()?; } + { + let target_id_mapping_file = get_target_id_mapping_file(&target_path); + let _file_lock = config.file_locks.write_lock(&target_id_mapping_file)?; + if let Ok(mut target_id_mapping) = BPlusTreeUpdate::::try_new(&target_id_mapping_file) { + if let Some(record) = target_id_mapping.query(&series_info_id) { + let new_record = record.copy_update_timestamp(); + let _ = target_id_mapping.update(&series_info_id, new_record); + } + }; + } + Ok(()) } // Reads the series info entry if exists, otherwise error -pub(crate) fn xtream_load_series_info(config: &Config, target_name: &str, series_id: u32) -> Result { - if let Some(storage_path) = xtream_get_storage_path(config, target_name) { - if let Some((info_path, idx_path)) = xtream_get_info_file_paths(&storage_path, XtreamCluster::Series) { - if info_path.exists() && idx_path.exists() { - return IndexedDocumentReader::::read_indexed_item(&info_path, &idx_path, series_id); +pub(crate) fn xtream_load_series_info(config: &Config, target_name: &str, series_id: u32) -> Option { + let target_path = get_target_storage_path(config, target_name)?; + let storage_path = xtream_get_storage_path(config, target_name)?; + + { + let target_id_mapping_file = get_target_id_mapping_file(&target_path); + let _file_lock = config.file_locks.read_lock(&target_id_mapping_file).map_err(|err| { + error!("Could not lock id mapping for target {target_name}: {}", err); + Error::new(ErrorKind::Other, format!("ID mapping load error for target {target_name}")) + }).ok()?; + let mut target_id_mapping = BPlusTreeQuery::::try_new(&target_id_mapping_file) + .map_err(|err| { + error!("Could not load id mapping for target {target_name}: {}", err); + Error::new(ErrorKind::Other, format!("ID mapping load error for target {target_name}")) + }).ok()?; + + if let Some(id_record) = target_id_mapping.query(&series_id) { + if id_record.is_expired() { + return None; } } } - Err(Error::new(ErrorKind::Other, format!("Failed to read series info for id {series_id} for {target_name}"))) -} -pub(crate) fn xtream_load_catchup_id_mapping(config: &Config, target_name: &str) -> Result, Error> { - if let Some(storage_path) = xtream_get_storage_path(config, target_name) { - let catchup_file = xtream_get_catchup_id_mapping_file_path(&storage_path); - return Ok(IdMapping::::new(&catchup_file)); - } - Err(Error::new(ErrorKind::Other, format!("Failed to load catchup id mapping {target_name}"))) -} + let (info_path, idx_path) = xtream_get_info_file_paths(&storage_path, XtreamCluster::Series)?; -pub(crate) fn write_and_get_xtream_series_info(config: &Config, target: &ConfigTarget, pli_series_info: &XtreamPlaylistItem, content: &str) -> Result { - if let Ok(mut doc) = serde_json::from_str::(content) { - if let Some(target_path) = get_target_storage_path(config, target.name.as_str()) { - let mut target_id_mapping = TargetIdMapping::new(&get_target_id_mapping_file(&target_path)); - if let Some(episodes) = doc.get_mut("episodes").and_then(|e| e.as_object_mut()) { - let options = XtreamMappingOptions::from_target_options(target.options.as_ref()); - for episode_list in episodes.values_mut() { - if let Some(entries) = episode_list.as_array_mut() { - for episode in entries.iter_mut().filter_map(|e| e.as_object_mut()) { - if let Some(episode_id) = episode.get("id").and_then(|id| id.as_str()) { - if let Ok(provider_id) = episode_id.parse::() { - let uuid = hash_string(&format!("{}/{}", pli_series_info.url, provider_id)); - let virtual_id = target_id_mapping.insert_entry(provider_id, uuid, &PlaylistItemType::Series, pli_series_info.virtual_id); - episode.insert("id".to_string(), Value::String(virtual_id.to_string())); - } - } - - if options.skip_series_direct_source { - episode.insert("direct_source".to_string(), Value::String(String::new())); - } - } - } + if info_path.exists() && idx_path.exists() { + { + let _file_lock = config.file_locks.read_lock(&info_path).map_err(|err| { + error!("Could not lock document {:?}: {}", info_path, err); + Error::new(ErrorKind::Other, format!("Document Reader error for target {target_name}")) + }).ok()?; + return match IndexedDocumentReader::::read_indexed_item(&info_path, &idx_path, series_id) { + Ok(content) => Some(content), + Err(err) => { + error!("Failed to read series info for id {series_id} for {target_name}: {}", err); + None } - - drop(target_id_mapping); - if let Ok(result) = serde_json::to_string(&doc) { - let _ = xtream_write_series_info(config, target.name.as_str(), pli_series_info.virtual_id, &result); - return Ok(result); - } - } - - - // if let Some(episodes) = doc.get_mut("episodes") { - // if let Some(episodes_map) = episodes.as_object_mut() { - // let options = XtreamMappingOptions::from_target_options(target.options.as_ref()); - // for (_season, episode_list) in episodes_map { - // // Iterate over items in the episode - // if let Some(entries) = episode_list.as_array_mut() { - // for entry in entries { - // if let Some(episode) = entry.as_object_mut() { - // if let Some(episode_id) = episode.get("id") { - // if let Ok(provider_id) = episode_id.as_str().unwrap().parse::() { - // let uuid = hash_string(format!("{}/{}", pli_series_info.url, provider_id).as_str()); - // let virtual_id = target_id_mapping.insert_entry(provider_id, uuid, &PlaylistItemType::Series, pli_series_info.virtual_id); - // episode.insert("id".to_string(), Value::String(virtual_id.to_string())); - // } - // } - // if options.skip_series_direct_source { - // episode.insert("direct_source".to_string(), Value::String(String::new())); - // } - // } - // } - // } - // } - // } - // drop(target_id_mapping); - // if let Ok(result) = serde_json::to_string(&doc) { - // let _ = xtream_repository::xtream_write_series_info(config, target.name.as_str(), pli_series_info.virtual_id, &result); - // return Ok(result); - // } - // } + }; } } - Err(Error::new(ErrorKind::Other, format!("Failed to get series info for id {}", pli_series_info.virtual_id))) + None +} + +pub(crate) fn write_and_get_xtream_series_info( + config: &Config, + target: &ConfigTarget, + pli_series_info: &XtreamPlaylistItem, + content: &str, +) -> Result { + let mut doc = serde_json::from_str::(content) + .map_err(|_| Error::new(ErrorKind::Other, "Failed to parse JSON content"))?; + + let target_path = get_target_storage_path(config, target.name.as_str()) + .ok_or_else(|| Error::new(ErrorKind::Other, format!("Could not find path for target {}", target.name)))?; + + let episodes = doc.get_mut("episodes") + .and_then(Value::as_object_mut) + .ok_or_else(|| Error::new(ErrorKind::Other, "No episodes found in content"))?; + + { + let target_id_mapping_file = get_target_id_mapping_file(&target_path); + let _file_lock = config.file_locks.write_lock(&target_id_mapping_file) + .map_err(|err| Error::new(ErrorKind::Other, format!("Could not load id mapping for target {} err:{}", target.name, err.to_string())))?; + let mut target_id_mapping = TargetIdMapping::new(&target_id_mapping_file); + let options = XtreamMappingOptions::from_target_options(target.options.as_ref()); + + for episode_list in episodes.values_mut().filter_map(Value::as_array_mut) { + for episode in episode_list.iter_mut().filter_map(Value::as_object_mut) { + if let Some(provider_id) = episode.get("id").and_then(Value::as_str).and_then(|id| id.parse::().ok()) { + let uuid = hash_string(&format!("{}/{}", pli_series_info.url, provider_id)); + let virtual_id = target_id_mapping.insert_entry(uuid, provider_id, &PlaylistItemType::SeriesEpisode, pli_series_info.virtual_id); + episode.insert("id".to_string(), Value::String(virtual_id.to_string())); + } + if options.skip_series_direct_source { + episode.insert("direct_source".to_string(), Value::String(String::new())); + } + } + } + + drop(target_id_mapping); + } + let result = serde_json::to_string(&doc) + .map_err(|_| Error::new(ErrorKind::Other, "Failed to serialize updated series info"))?; + xtream_write_series_info(config, target.name.as_str(), pli_series_info.virtual_id, &result).ok(); + + Ok(result) } \ No newline at end of file diff --git a/src/utils/file_lock_manager.rs b/src/utils/file_lock_manager.rs new file mode 100644 index 000000000..6e16675e3 --- /dev/null +++ b/src/utils/file_lock_manager.rs @@ -0,0 +1,106 @@ +use std::collections::HashMap; +use std::sync::{Arc, Mutex, RwLock, RwLockReadGuard, RwLockWriteGuard}; +use std::{fmt, io}; +use std::path::{Path, PathBuf}; + +#[derive(Clone)] +pub(crate) struct FileLockManager { + locks: Arc>>>>, +} + +impl FileLockManager { + pub(crate) fn new() -> Self { + Self { + locks: Arc::new(Mutex::new(HashMap::new())), + } + } + + // Acquires a read lock for the specified file and returns a FileReadGuard. + pub(crate) fn read_lock(&self, path: &Path) -> io::Result { + let file_lock = self.get_or_create_lock(path)?; + let guard = file_lock.read().map_err(|_| { + io::Error::new(io::ErrorKind::Other, "Failed to acquire read lock") + })?; + // Clone the Arc to avoid moving `file_lock` out, as it is still borrowed by `guard` + Ok(FileReadGuard::new(Arc::clone(&file_lock), guard)) + } + + // Acquires a write lock for the specified file and returns a FileWriteGuard. + pub(crate) fn write_lock(&self, path: &Path) -> io::Result { + let file_lock = self.get_or_create_lock(path)?; + let guard = file_lock.write().map_err(|_| { + io::Error::new(io::ErrorKind::Other, "Failed to acquire write lock") + })?; + // Clone the Arc to avoid moving `file_lock` out, as it is still borrowed by `guard` + Ok(FileWriteGuard::new(Arc::clone(&file_lock), guard)) + } + + // Helper function: retrieves or creates a lock for a file. + fn get_or_create_lock(&self, path: &Path) -> io::Result>> { + let mut locks = self.locks.lock().map_err(|_| { + io::Error::new(io::ErrorKind::Other, "Failed to acquire lock on lock manager") + })?; + + if let Some(lock) = locks.get(path) { + return Ok(lock.clone()); + } + + let file_lock = Arc::new(RwLock::new(())); + locks.insert(path.to_path_buf(), file_lock.clone()); + Ok(file_lock) + } +} + +impl Default for FileLockManager { + fn default() -> Self { + FileLockManager::new() + } +} + +impl fmt::Debug for FileLockManager { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + // Acquire the lock to safely access the HashMap + let locks = self.locks.lock().unwrap(); + // Format the paths in the HashMap for debug purposes + let keys: Vec<_> = locks.keys().collect(); + f.debug_struct("FileLockManager") + .field("locks", &keys) + .finish() + } +} + +// Define FileReadGuard to hold both the lock reference and the actual read guard. +#[allow(dead_code)] +pub(crate) struct FileReadGuard { + lock: Arc>, + guard: RwLockReadGuard<'static, ()>, +} + +impl FileReadGuard { + pub fn new(lock: Arc>, guard: RwLockReadGuard<'_, ()>) -> Self { + // Convert the lifetime of `guard` to 'static by transmuting. + let static_guard: RwLockReadGuard<'static, ()> = unsafe { std::mem::transmute(guard) }; + FileReadGuard { + lock, + guard: static_guard, + } + } +} + +// Define FileWriteGuard to hold both the lock reference and the actual write guard. +#[allow(dead_code)] +pub(crate) struct FileWriteGuard { + lock: Arc>, + guard: RwLockWriteGuard<'static, ()>, +} + +impl FileWriteGuard { + pub fn new(lock: Arc>, guard: RwLockWriteGuard<'_, ()>) -> Self { + // Convert the lifetime of `guard` to 'static by transmuting. + let static_guard: RwLockWriteGuard<'static, ()> = unsafe { std::mem::transmute(guard) }; + FileWriteGuard { + lock, + guard: static_guard, + } + } +} diff --git a/src/utils/mod.rs b/src/utils/mod.rs index 008bd750c..be648a6f0 100644 --- a/src/utils/mod.rs +++ b/src/utils/mod.rs @@ -5,4 +5,5 @@ pub (crate) mod string_utils; pub (crate) mod json_utils; pub (crate) mod config_reader; pub (crate) mod default_utils; -pub (crate) mod multi_file_reader; \ No newline at end of file +pub (crate) mod multi_file_reader; +pub (crate) mod file_lock_manager; \ No newline at end of file