diff --git a/Cargo.lock b/Cargo.lock index 90f55fac6..f48b289dc 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2406,6 +2406,15 @@ version = "0.1.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "112b39cec0b298b6c1999fee3e31427f74f676e4cb9879ed1a121b43661a4154" +[[package]] +name = "lz4_flex" +version = "0.12.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ab6473172471198271ff72e9379150e9dfd70d8e533e0752a27e515b48dd375e" +dependencies = [ + "twox-hash", +] + [[package]] name = "matchit" version = "0.8.4" @@ -3587,15 +3596,6 @@ version = "1.0.22" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "b39cdef0fa800fc44525c84ccb54a029961a8215f9619753635a9c0d2538d46d" -[[package]] -name = "ruzstd" -version = "0.8.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e5ff0cc5e135c8870a775d3320910cd9b564ec036b4dc0b8741629020be63f01" -dependencies = [ - "twox-hash", -] - [[package]] name = "ryu" version = "1.0.20" @@ -3836,6 +3836,7 @@ dependencies = [ "indexmap", "js-sys", "log", + "lz4_flex", "path-clean", "pest", "pest_derive", @@ -4433,7 +4434,6 @@ dependencies = [ "rpassword", "rphonetic", "rust-argon2", - "ruzstd", "serde", "serde_json", "serde_yaml", diff --git a/backend/Cargo.toml b/backend/Cargo.toml index 1bd25e2bb..2553137a8 100644 --- a/backend/Cargo.toml +++ b/backend/Cargo.toml @@ -50,7 +50,6 @@ tokio = { version = "1.48", features = ["rt-multi-thread", "parking_lot", "fs"] #tracing-subscriber = { version = "0.3", features = ["fmt", "env-filter"] } tokio-util = { version = "0.7", features = ["io", "io-util"] } tempfile = "3.23" -ruzstd = "0.8" filetime = "0.2" strsim = "0.11" rphonetic = "3.0" diff --git a/backend/src/main.rs b/backend/src/main.rs index 2f410e230..415100d31 100644 --- a/backend/src/main.rs +++ b/backend/src/main.rs @@ -111,7 +111,7 @@ async fn main() { match (args.db_content_type.as_ref(), args.db_file_name.as_ref()) { (Some(dbt), Some(dbf)) => db_viewer(dbf, dbt), - (Some(_) | None, None) | (None, Some(_))=> info!("You need to define database content type and filename!"), + (Some(_) | None, None) | (None, Some(_))=> eprintln!("You need to define database content type and filename!"), } if args.genpwd { diff --git a/backend/src/processing/processor/xtream_series.rs b/backend/src/processing/processor/xtream_series.rs index b35d4c749..b75cf10c4 100644 --- a/backend/src/processing/processor/xtream_series.rs +++ b/backend/src/processing/processor/xtream_series.rs @@ -1,23 +1,22 @@ -use crate::model::{AppConfig, ConfigTarget}; use crate::model::FetchedPlaylist; +use crate::model::{AppConfig, ConfigTarget}; use crate::processing::parser::xtream::parse_xtream_series_info; use crate::processing::processor::playlist::ProcessingPipe; -use crate::processing::processor::xtream::{playlist_resolve_download_playlist_item}; -use crate::processing::processor::{create_resolve_options_function_for_xtream_target}; +use crate::processing::processor::xtream::playlist_resolve_download_playlist_item; +use crate::processing::processor::create_resolve_options_function_for_xtream_target; +use crate::repository::storage::get_input_storage_path; +use crate::repository::xtream_repository::persists_input_series_info; use log::{error, info, log_enabled, Level}; use shared::error::TuliproxError; use shared::model::{InputType, PlaylistEntry, SeriesStreamProperties, StreamProperties, XtreamSeriesInfo}; use shared::model::{PlaylistGroup, PlaylistItemType, XtreamCluster}; use std::sync::Arc; use std::time::Instant; -use crate::repository::storage::get_input_storage_path; -use crate::repository::xtream_repository::{persists_input_series_info}; create_resolve_options_function_for_xtream_target!(series); async fn playlist_resolve_series_info(app_config: &Arc, client: &reqwest::Client, errors: &mut Vec, fpl: &mut FetchedPlaylist<'_>, resolve_delay: u16) -> Vec { - let input = fpl.input; let working_dir = &app_config.config.load().working_dir; let storage_path = match get_input_storage_path(&input.name, working_dir) { @@ -43,38 +42,43 @@ async fn playlist_resolve_series_info(app_config: &Arc, client: &reqw for plg in &mut fpl.playlistgroups { let mut group_series = vec![]; for pli in &mut plg.channels { - processed_series_info_count += 1; if pli.header.xtream_cluster != XtreamCluster::Series || pli.header.item_type != PlaylistItemType::SeriesInfo || pli.has_details() { continue; } let Some(provider_id) = pli.get_provider_id() else { continue; }; + processed_series_info_count += 1; let (group, series_name) = { let header = &pli.header; (header.group.clone(), if header.name.is_empty() { header.title.clone() } else { header.name.clone() }) }; if provider_id != 0 { if let Some(content) = playlist_resolve_download_playlist_item(client, pli, fpl.input, errors, resolve_delay, XtreamCluster::Series).await { - if let Ok(info) = serde_json::from_str::(&content) { - let series_stream_props = SeriesStreamProperties::from_info(&info, pli); - // the input db needs to be updated - let _ = persists_input_series_info(app_config, &storage_path, pli.header.xtream_cluster, &input.name, provider_id, &series_stream_props).await; - // extract episodes from info - if let Some(episodes) = parse_xtream_series_info(&pli.get_uuid(), &series_stream_props, &group, &series_name, input) { - group_series.extend(episodes.into_iter()); + match serde_json::from_str::(&content) { + Ok(info) => { + let series_stream_props = SeriesStreamProperties::from_info(&info, pli); + // the input db needs to be updated + let _ = persists_input_series_info(app_config, &storage_path, pli.header.xtream_cluster, &input.name, provider_id, &series_stream_props).await; + // extract episodes from info + if let Some(episodes) = parse_xtream_series_info(&pli.get_uuid(), &series_stream_props, &group, &series_name, input) { + group_series.extend(episodes.into_iter()); + } + // Update in-memory playlist items with the newly fetched vod info. + // This makes the data available for later processing steps like STRM export. + pli.header.additional_properties = Some(StreamProperties::Series(Box::new(series_stream_props))); + } + Err(err) => { + error!("Failed to parse series info for provider_id {provider_id}: {err}"); } - // Update in-memory playlist items with the newly fetched vod info. - // This makes the data available for later processing steps like STRM export. - pli.header.additional_properties = Some(StreamProperties::Series(Box::new(series_stream_props))); } } } if log_enabled!(Level::Info) && last_log_time.elapsed().as_secs() >= 30 { - info!("resolved {processed_series_info_count}/{series_info_count} series info"); - last_log_time = Instant::now(); - } + info!("resolved {processed_series_info_count}/{series_info_count} series info"); + last_log_time = Instant::now(); + } } if !group_series.is_empty() { result.push(PlaylistGroup { diff --git a/backend/src/processing/processor/xtream_vod.rs b/backend/src/processing/processor/xtream_vod.rs index e56c39a59..b50460da0 100644 --- a/backend/src/processing/processor/xtream_vod.rs +++ b/backend/src/processing/processor/xtream_vod.rs @@ -1,15 +1,15 @@ use crate::model::FetchedPlaylist; use crate::model::{AppConfig, ConfigTarget}; -use crate::processing::processor::xtream::playlist_resolve_download_playlist_item; use crate::processing::processor::create_resolve_options_function_for_xtream_target; +use crate::processing::processor::xtream::playlist_resolve_download_playlist_item; use crate::repository::storage::get_input_storage_path; +use crate::repository::xtream_repository::persist_input_vod_info; use log::{error, info, log_enabled, Level}; use shared::error::TuliproxError; use shared::model::{InputType, PlaylistEntry, StreamProperties, VideoStreamProperties, XtreamVideoInfo}; use shared::model::{PlaylistItemType, XtreamCluster}; use std::sync::Arc; use std::time::Instant; -use crate::repository::xtream_repository::{persist_input_vod_info}; create_resolve_options_function_for_xtream_target!(vod); @@ -41,22 +41,27 @@ pub async fn playlist_resolve_vod(app_config: &Arc, client: &reqwest: for plg in &mut fpl.playlistgroups { for pli in &mut plg.channels { - processed_vod_info_count += 1; if pli.header.xtream_cluster != XtreamCluster::Video || pli.header.item_type != PlaylistItemType::Video || pli.has_details() { continue; } let Some(provider_id) = pli.get_provider_id() else { continue; }; + processed_vod_info_count += 1; if provider_id != 0 { if let Some(content) = playlist_resolve_download_playlist_item(client, pli, fpl.input, errors, resolve_delay, XtreamCluster::Video).await { - if let Ok(info) = serde_json::from_str::(&content) { - let video_stream_props = VideoStreamProperties::from_info(&info, pli); - if let Err(err) = persist_input_vod_info(app_config, &storage_path, pli.header.xtream_cluster, &input.name, provider_id, &video_stream_props).await { - error!("Failed to persist VOD info for provider_id {provider_id}: {err}"); + match serde_json::from_str::(&content) { + Ok(info) => { + let video_stream_props = VideoStreamProperties::from_info(&info, pli); + if let Err(err) = persist_input_vod_info(app_config, &storage_path, pli.header.xtream_cluster, &input.name, provider_id, &video_stream_props).await { + error!("Failed to persist VOD info for provider_id {provider_id}: {err}"); + } + // This makes the data available for subsequent processing steps like STRM export. + pli.header.additional_properties = Some(StreamProperties::Video(Box::new(video_stream_props))); + } + Err(err) => { + error!("Failed to parse video info for provider_id {provider_id}: {err}"); } - // This makes the data available for subsequent processing steps like STRM export. - pli.header.additional_properties = Some(StreamProperties::Video(Box::new(video_stream_props))); } } } diff --git a/backend/src/repository/bplustree.rs b/backend/src/repository/bplustree.rs index fdbf2ff1c..99d24176b 100644 --- a/backend/src/repository/bplustree.rs +++ b/backend/src/repository/bplustree.rs @@ -1187,6 +1187,10 @@ where self.leaf_values = values; self.leaf_idx = 0; } + Err(err) => { + error!("BPlusTreeDiskIterator Failed to read next entry: {err}"); + return None + } _ => return None, } } @@ -1319,7 +1323,7 @@ where let child_idx = get_entry_index_upper_bound::(&node.keys, key); if let Some(mut pters) = pointers { if let Some(&child_offset) = pters.get(child_idx) { - // Recurse to get new child offset + // Recurse to get a new child offset let new_child_offset = self.update_recursive(child_offset, key, value)?; // Update current node's pointers @@ -1354,6 +1358,108 @@ where Ok(new_root_offset) } + /// Update multiple items in batch. This is more efficient than calling `update()` multiple times + /// as it performs all updates and then commits the final root offset once. + /// returns The final root offset after all updates, or an error if any update fails + pub fn update_batch(&mut self, items: &[(&K, &V)]) -> io::Result { + if items.is_empty() { + return Ok(self.root_offset); + } + + let mut current_root = self.root_offset; + + // Perform all updates sequentially + for (key, value) in items { + current_root = self.update_recursive(current_root, key, value)?; + } + + // Atomic Header Swap - only once at the end + self.file.seek(SeekFrom::Start(ROOT_OFFSET_POS))?; + self.file.write_all(¤t_root.to_le_bytes())?; + self.file.flush()?; + self.file.sync_all()?; + + self.root_offset = current_root; + Ok(current_root) + } + + /// Insert or update multiple items in batch (upsert). If a key exists, it will be updated; + /// if it doesn't exist, it will be inserted. This is more efficient than calling `update()` + /// or `insert()` multiple times as it loads the tree once, performs all operations, and saves once. + /// returns The final root offset after all upserts, or an error if any operation fails + pub fn upsert_batch(&mut self, items: &[(&K, &V)]) -> io::Result { + if items.is_empty() { + return Ok(self.root_offset); + } + + // Get the filepath from the file handle + // We need to load the tree, so we'll use a temporary approach + // Load the current tree state into memory + let mut tree = { + // Create a reader from our file + let mut reader = utils::file_reader(&mut self.file); + + // Read header to get root offset + let mut header = [0u8; 16]; + reader.seek(SeekFrom::Start(0))?; + reader.read_exact(&mut header)?; + + if &header[0..4] != MAGIC { + return Err(io::Error::new(io::ErrorKind::InvalidData, "Invalid magic number")); + } + let version = u32::from_le_bytes(header[4..8].try_into().map_err(|_| io::Error::new(io::ErrorKind::InvalidData, "Invalid version slice"))?); + if version != STORAGE_VERSION { + return Err(io::Error::new(io::ErrorKind::InvalidData, format!("Unsupported storage version: {version}"))); + } + + // Load the tree from current root offset + let mut buffer = vec![0u8; BLOCK_SIZE]; + let (root, _) = BPlusTreeNode::::deserialize_from_block(&mut reader, &mut buffer, self.root_offset, true)?; + + let (inner_order, leaf_order) = calc_order::(); + BPlusTree { + root, + inner_order, + leaf_order, + dirty: false, + } + }; + + // Perform all inserts/updates + for (key, value) in items { + tree.insert((*key).clone(), (*value).clone()); + } + + // Mark as dirty and save using store_internal (we already hold the lock) + tree.dirty = true; + + // We need to save to a temp file and then rename, but we can't easily get the filepath + // Instead, we'll serialize directly to our file handle + // Seek to beginning and write + self.file.seek(SeekFrom::Start(0))?; + self.file.set_len(0)?; // Truncate the file + + let mut buffer = vec![0u8; BLOCK_SIZE]; + + // Write header block + let mut header = [0u8; BLOCK_SIZE]; + header[0..4].copy_from_slice(MAGIC); + header[4..8].copy_from_slice(&STORAGE_VERSION.to_le_bytes()); + header[8..16].copy_from_slice(&HEADER_SIZE.to_le_bytes()); + self.file.write_all(&header)?; + + // Serialize tree + tree.root.serialize_breadth_first(&mut self.file, &mut buffer, HEADER_SIZE)?; + + self.file.flush()?; + self.file.sync_all()?; + + // Update our root offset + self.root_offset = HEADER_SIZE; + + Ok(self.root_offset) + } + /// Garbage Collection: Compacts the file by rewriting only live blocks sequentially. pub fn compact(&mut self, filepath: &Path) -> io::Result<()> { // 1. Reload the current tree fully from the live root @@ -1830,4 +1936,337 @@ mod tests { Ok(()) } + + #[test] + fn update_batch_basic_test() -> io::Result<()> { + let tempdir = tempfile::tempdir()?; + let filepath = tempdir.path().join("tree_update_batch.bin"); + let mut tree = BPlusTree::::new(); + + // Create initial tree + for i in 0u32..50 { + tree.insert(i, Record { + id: i, + data: format!("Initial {i}") + }); + } + tree.store(&filepath)?; + drop(tree); + + // Test batch update + let mut tree_update = BPlusTreeUpdate::::try_new(&filepath)?; + + // Prepare batch updates + let updates: Vec<(u32, Record)> = (0u32..50) + .filter(|i| i % 5 == 0) + .map(|i| (i, Record { id: i, data: format!("BatchUpdated {i}") })) + .collect(); + + let update_refs: Vec<(&u32, &Record)> = updates.iter() + .map(|(k, v)| (k, v)) + .collect(); + + tree_update.update_batch(&update_refs)?; + drop(tree_update); + + // Verify all updates + let mut tree_query = BPlusTreeQuery::::try_new(&filepath)?; + for i in 0u32..50 { + let val = tree_query.query(&i).expect("Should find key"); + if i % 5 == 0 { + assert_eq!(val.data, format!("BatchUpdated {i}"), "Batch updated key {i} should have new value"); + } else { + assert_eq!(val.data, format!("Initial {i}"), "Non-updated key {i} should have original value"); + } + } + + Ok(()) + } + + #[test] + fn update_batch_empty_test() -> io::Result<()> { + let tempdir = tempfile::tempdir()?; + let filepath = tempdir.path().join("tree_update_batch_empty.bin"); + let mut tree = BPlusTree::::new(); + + // Create initial tree + for i in 0u32..10 { + tree.insert(i, Record { + id: i, + data: format!("Initial {i}") + }); + } + tree.store(&filepath)?; + drop(tree); + + let mut tree_update = BPlusTreeUpdate::::try_new(&filepath)?; + let initial_root = tree_update.root_offset; + + // Test empty batch - should be no-op + let empty_batch: Vec<(&u32, &Record)> = vec![]; + let result = tree_update.update_batch(&empty_batch)?; + + assert_eq!(result, initial_root, "Empty batch should not change root offset"); + + // Verify data unchanged + for i in 0u32..10 { + let val = tree_update.query(&i).expect("Should find key"); + assert_eq!(val.data, format!("Initial {i}")); + } + + Ok(()) + } + + #[test] + fn update_batch_large_test() -> io::Result<()> { + let tempdir = tempfile::tempdir()?; + let filepath = tempdir.path().join("tree_update_batch_large.bin"); + let mut tree = BPlusTree::::new(); + + let test_size = 200u32; + + // Create initial tree + for i in 0..test_size { + tree.insert(i, Record { + id: i, + data: format!("Initial {i}") + }); + } + tree.store(&filepath)?; + drop(tree); + + let mut tree_update = BPlusTreeUpdate::::try_new(&filepath)?; + + // Prepare large batch update (every other item) + let updates: Vec<(u32, Record)> = (0..test_size) + .filter(|i| i % 2 == 0) + .map(|i| (i, Record { id: i, data: format!("BatchUpdated {i}") })) + .collect(); + + let update_refs: Vec<(&u32, &Record)> = updates.iter() + .map(|(k, v)| (k, v)) + .collect(); + + // Perform batch update + tree_update.update_batch(&update_refs)?; + drop(tree_update); + + // Verify all updates via iterator + let reloaded_tree = BPlusTree::::load(&filepath)?; + let mut count = 0; + for (key, value) in &reloaded_tree { + if *key % 2 == 0 { + assert_eq!(value.data, format!("BatchUpdated {key}"), "Even keys should be batch updated"); + } else { + assert_eq!(value.data, format!("Initial {key}"), "Odd keys should remain unchanged"); + } + count += 1; + } + assert_eq!(count, test_size, "Should have all entries"); + + Ok(()) + } + + #[test] + fn update_batch_with_compaction_test() -> io::Result<()> { + let tempdir = tempfile::tempdir()?; + let filepath = tempdir.path().join("tree_update_batch_compact.bin"); + let mut tree = BPlusTree::::new(); + + // Create initial tree with larger data + let large_data = "x".repeat(1000); + for i in 0u32..100 { + tree.insert(i, Record { + id: i, + data: large_data.clone() + }); + } + tree.store(&filepath)?; + drop(tree); + + let mut tree_update = BPlusTreeUpdate::::try_new(&filepath)?; + + // Batch update with smaller data + let small_data = "y".repeat(50); + let updates: Vec<(u32, Record)> = (0u32..100) + .map(|i| (i, Record { id: i, data: small_data.clone() })) + .collect(); + + let update_refs: Vec<(&u32, &Record)> = updates.iter() + .map(|(k, v)| (k, v)) + .collect(); + + tree_update.update_batch(&update_refs)?; + + let size_before_compact = std::fs::metadata(&filepath)?.len(); + + // Compact to reclaim space + tree_update.compact(&filepath)?; + + let size_after_compact = std::fs::metadata(&filepath)?.len(); + assert!(size_after_compact < size_before_compact, + "Compaction should reduce file size after batch update"); + + // Verify all data is correct after compaction + drop(tree_update); + let mut tree_query = BPlusTreeQuery::::try_new(&filepath)?; + for i in 0u32..100 { + let val = tree_query.query(&i).expect("Should find key after compaction"); + assert_eq!(val.data, small_data, "Data should be updated after compaction"); + } + + Ok(()) + } + + #[test] + fn upsert_batch_mixed_test() -> io::Result<()> { + let tempdir = tempfile::tempdir()?; + let filepath = tempdir.path().join("tree_upsert_batch_mixed.bin"); + let mut tree = BPlusTree::::new(); + + // Create initial tree with keys 0-49 + for i in 0u32..50 { + tree.insert(i, Record { + id: i, + data: format!("Initial {i}") + }); + } + tree.store(&filepath)?; + drop(tree); + + let mut tree_update = BPlusTreeUpdate::::try_new(&filepath)?; + + // Prepare upsert batch: update existing keys 0-24, insert new keys 50-74 + let mut updates: Vec<(u32, Record)> = Vec::new(); + + // Updates to existing keys + for i in 0u32..25 { + updates.push((i, Record { id: i, data: format!("Updated {i}") })); + } + + // Inserts for new keys + for i in 50u32..75 { + updates.push((i, Record { id: i, data: format!("Inserted {i}") })); + } + + let update_refs: Vec<(&u32, &Record)> = updates.iter() + .map(|(k, v)| (k, v)) + .collect(); + + tree_update.upsert_batch(&update_refs)?; + drop(tree_update); + + // Verify all 75 entries exist with correct values + let mut tree_query = BPlusTreeQuery::::try_new(&filepath)?; + + // Check updated keys (0-24) + for i in 0u32..25 { + let val = tree_query.query(&i).expect("Should find updated key"); + assert_eq!(val.data, format!("Updated {i}"), "Key {i} should be updated"); + } + + // Check unchanged keys (25-49) + for i in 25u32..50 { + let val = tree_query.query(&i).expect("Should find unchanged key"); + assert_eq!(val.data, format!("Initial {i}"), "Key {i} should remain unchanged"); + } + + // Check inserted keys (50-74) + for i in 50u32..75 { + let val = tree_query.query(&i).expect("Should find inserted key"); + assert_eq!(val.data, format!("Inserted {i}"), "Key {i} should be inserted"); + } + + Ok(()) + } + + #[test] + fn upsert_batch_all_new_test() -> io::Result<()> { + let tempdir = tempfile::tempdir()?; + let filepath = tempdir.path().join("tree_upsert_batch_new.bin"); + let mut tree = BPlusTree::::new(); + + // Create initial tree with unrelated keys + for i in 0u32..10 { + tree.insert(i, Record { + id: i, + data: format!("Initial {i}") + }); + } + tree.store(&filepath)?; + drop(tree); + + let mut tree_update = BPlusTreeUpdate::::try_new(&filepath)?; + + // Upsert all new keys (100-149) + let updates: Vec<(u32, Record)> = (100u32..150) + .map(|i| (i, Record { id: i, data: format!("New {i}") })) + .collect(); + + let update_refs: Vec<(&u32, &Record)> = updates.iter() + .map(|(k, v)| (k, v)) + .collect(); + + tree_update.upsert_batch(&update_refs)?; + drop(tree_update); + + // Verify all keys exist + let mut tree_query = BPlusTreeQuery::::try_new(&filepath)?; + + // Original keys should still exist + for i in 0u32..10 { + let val = tree_query.query(&i).expect("Should find original key"); + assert_eq!(val.data, format!("Initial {i}")); + } + + // New keys should be inserted + for i in 100u32..150 { + let val = tree_query.query(&i).expect("Should find new key"); + assert_eq!(val.data, format!("New {i}")); + } + + Ok(()) + } + + #[test] + fn upsert_batch_all_existing_test() -> io::Result<()> { + let tempdir = tempfile::tempdir()?; + let filepath = tempdir.path().join("tree_upsert_batch_existing.bin"); + let mut tree = BPlusTree::::new(); + + // Create initial tree + for i in 0u32..100 { + tree.insert(i, Record { + id: i, + data: format!("Initial {i}") + }); + } + tree.store(&filepath)?; + drop(tree); + + let mut tree_update = BPlusTreeUpdate::::try_new(&filepath)?; + + // Upsert all existing keys (should behave like update) + let updates: Vec<(u32, Record)> = (0u32..100) + .map(|i| (i, Record { id: i, data: format!("Updated {i}") })) + .collect(); + + let update_refs: Vec<(&u32, &Record)> = updates.iter() + .map(|(k, v)| (k, v)) + .collect(); + + tree_update.upsert_batch(&update_refs)?; + drop(tree_update); + + // Verify all values were updated + let reloaded_tree = BPlusTree::::load(&filepath)?; + let mut count = 0; + for (key, value) in &reloaded_tree { + assert_eq!(value.data, format!("Updated {key}"), "All keys should be updated"); + count += 1; + } + assert_eq!(count, 100, "Should have exactly 100 entries"); + + Ok(()) + } } diff --git a/backend/src/repository/xtream_repository.rs b/backend/src/repository/xtream_repository.rs index 61011972c..167bd4af5 100644 --- a/backend/src/repository/xtream_repository.rs +++ b/backend/src/repository/xtream_repository.rs @@ -126,7 +126,7 @@ async fn write_playlists_to_file( // Ok(()) // } -pub async fn write_playlist_item_to_file( +pub async fn write_playlist_item_update( app_config: &Arc, target_name: &str, pli: &XtreamPlaylistItem, @@ -144,11 +144,37 @@ pub async fn write_playlist_item_to_file( // This case should rarely happen as the file is usually pre-created, but for safety: return Err(cant_write_result!(&xtream_path, "BPlusTree file not found for append")); }; + tree.update(&pli.virtual_id, pli.clone()).map_err(|err| cant_write_result!(&xtream_path, err))?; } Ok(()) } +pub async fn write_playlist_batch_item_upsert( + app_config: &Arc, + target_name: &str, + xtream_cluster: XtreamCluster, + pli_list: &[XtreamPlaylistItem], +) -> Result<(), TuliproxError> { + let storage_path = { + let config = app_config.config.load(); + ensure_xtream_storage_path(&config, target_name)? + }; + let xtream_path = xtream_get_file_path(&storage_path, xtream_cluster); + { + let _file_lock = app_config.file_locks.write_lock(&xtream_path).await; + let mut tree = if xtream_path.exists() { + BPlusTreeUpdate::try_new(&xtream_path).map_err(|err| cant_write_result!(&xtream_path, err))? + } else { + // This case should rarely happen as the file is usually pre-created, but for safety: + return Err(cant_write_result!(&xtream_path, "BPlusTree file not found for append")); + }; + + let batch: Vec<(&u32, &XtreamPlaylistItem)> = pli_list.iter().map(|pli| (&pli.virtual_id, pli)).collect(); + tree.upsert_batch(&batch).map_err(|err| cant_write_result!(&xtream_path, err))?; + } + Ok(()) +} fn get_map_item_as_str(map: &serde_json::Map, key: &str) -> Option { if let Some(value) = map.get(key) { @@ -746,4 +772,5 @@ pub async fn persists_input_series_info(app_config: &Arc, storage_pat cluster: XtreamCluster, input_name: &str, provider_id: u32, props: &SeriesStreamProperties) -> Result<(), Error> { persist_input_info(app_config, storage_path, cluster, input_name, provider_id, StreamProperties::Series(Box::new(props.clone()))).await -} \ No newline at end of file +} + diff --git a/backend/src/utils/db_viewer.rs b/backend/src/utils/db_viewer.rs index df5b68042..8869e6248 100644 --- a/backend/src/utils/db_viewer.rs +++ b/backend/src/utils/db_viewer.rs @@ -1,5 +1,7 @@ use std::io::Write; use std::path::PathBuf; +use env_logger::{Builder, Target}; +use log::{error, warn, LevelFilter}; use serde::{Deserialize, Serialize}; use shared::model::{M3uPlaylistItem, XtreamPlaylistItem}; use crate::repository::bplustree::{BPlusTreeDiskIterator, BPlusTreeQuery}; @@ -8,17 +10,23 @@ pub fn db_viewer(filename: &str, content_type: &str) { let path = match PathBuf::from(filename).canonicalize() { Ok(p) => p, Err(err) => { - println!("File does not exist! {err}"); - let _ = std::io::stdout().flush(); + eprintln!("File does not exist! {err}"); + let _ = std::io::stderr().flush(); std::process::exit(1); } }; if !path.exists() { - println!("File does not exist! {}", path.display()); - let _ = std::io::stdout().flush(); + eprintln!("File does not exist! {}", path.display()); + let _ = std::io::stderr().flush(); std::process::exit(1); } + + let mut log_builder = Builder::from_default_env(); + log_builder.target(Target::Stderr); + log_builder.filter_level(LevelFilter::Info); + log_builder.init(); + match content_type { "xtream" => { if let Ok(mut query) = BPlusTreeQuery::::try_new(&path) { @@ -32,9 +40,10 @@ pub fn db_viewer(filename: &str, content_type: &str) { print_json_from_iter(iterator); } } - _ => println!("Allowed content types are: [m3u, xtream]"), + _ => warn!("Allowed content types are: [m3u, xtream]"), } let _ = std::io::stdout().flush(); + let _ = std::io::stderr().flush(); std::process::exit(1); } @@ -52,7 +61,7 @@ where P: Serialize + for<'de> Deserialize<'de> + Clone println!("{text}"); first = false; }, - Err(err) => eprintln!("Failed: {err}"), + Err(err) => error!("Failed: {err}"), } } println!("]"); diff --git a/backend/src/utils/network/xtream.rs b/backend/src/utils/network/xtream.rs index 68073abf5..43ccb6050 100644 --- a/backend/src/utils/network/xtream.rs +++ b/backend/src/utils/network/xtream.rs @@ -1,18 +1,20 @@ use crate::api::model::AppState; use crate::messaging::send_message; -use crate::model::{is_input_expired, xtream_mapping_option_from_target_options, Config, ConfigInput, ConfigTarget, XtreamLoginInfo}; +use crate::model::{is_input_expired, xtream_mapping_option_from_target_options, Config, ConfigInput, ConfigTarget, XtreamLoginInfo, XtreamTargetOutput}; use crate::model::{InputSource, ProxyUserCredentials}; use crate::processing::parser::xtream; use crate::processing::parser::xtream::parse_xtream_series_info; use crate::repository::playlist_repository::{get_target_id_mapping, rewrite_provider_series_info_episode_virtual_id, ProviderEpisodeKey}; use crate::repository::storage::{get_input_storage_path, get_target_storage_path}; use crate::repository::target_id_mapping::VirtualIdRecord; -use crate::repository::xtream_repository::{persist_input_vod_info, persists_input_series_info, write_playlist_item_to_file}; +use crate::repository::xtream_repository::{persist_input_vod_info, persists_input_series_info, write_playlist_batch_item_upsert, write_playlist_item_update}; use crate::utils::request; use chrono::{DateTime, Utc}; use log::{error, info, warn}; use shared::error::{str_to_io_error, to_io_error, TuliproxError}; -use shared::model::{MsgKind, PlaylistEntry, PlaylistGroup, ProxyUserStatus, SeriesStreamProperties, StreamProperties, VideoStreamProperties, XtreamCluster, XtreamPlaylistItem, XtreamSeriesInfo, XtreamVideoInfo}; +use shared::model::{MsgKind, PlaylistEntry, PlaylistGroup, ProxyUserStatus, SeriesStreamProperties, + StreamProperties, VideoStreamProperties, XtreamCluster, + XtreamPlaylistItem, XtreamSeriesInfo, XtreamVideoInfo}; use shared::utils::{extract_extension_from_url, get_i64_from_serde_value, get_string_from_serde_value, sanitize_sensitive_info}; use std::collections::HashMap; use std::io::Error; @@ -84,26 +86,33 @@ pub async fn get_xtream_stream_info(client: &reqwest::Client, XtreamCluster::Video => { let working_dir = &app_config.config.load().working_dir; if let Ok(storage_path) = get_input_storage_path(&input.name, working_dir) { - if let Ok(info) = serde_json::from_str::(&content) { - // parse downloaded info into StreamProperties - let video_stream_props = VideoStreamProperties::from_info(&info, pli); + match serde_json::from_str::(&content) { + Ok(info) => { + // parse downloaded info into StreamProperties + let video_stream_props = VideoStreamProperties::from_info(&info, pli); - // persist input info - if let Err(err) = persist_input_vod_info(&app_state.app_config, &storage_path, cluster, &input.name, provider_id, &video_stream_props).await { - error!("Failed to persist video stream for input {}: {err}", &input.name); - } - - // update target playlist - let mut vod_pli = pli.clone(); - vod_pli.additional_properties = Some(StreamProperties::Video(Box::new(video_stream_props))); - - if let Err(err) = write_playlist_item_to_file(app_config, &target.name, &vod_pli).await { - error!("Failed to persist video stream: {err}"); - } - - if target.use_memory_cache { - app_state.playlists.update_playlist_items(target, vec![&vod_pli]).await; + // persist input info + if let Err(err) = persist_input_vod_info(&app_state.app_config, &storage_path, cluster, &input.name, provider_id, &video_stream_props).await { + error!("Failed to persist video stream for input {}: {err}", &input.name); + } + + // update target playlist + let mut vod_pli = pli.clone(); + vod_pli.additional_properties = Some(StreamProperties::Video(Box::new(video_stream_props))); + + if let Err(err) = write_playlist_item_update(app_config, &target.name, &vod_pli).await { + error!("Failed to persist video stream: {err}"); + } + + if target.use_memory_cache { + app_state.playlists.update_playlist_items(target, vec![&vod_pli]).await; + } + + if let Some(value) = xtream_resolve_stream_info(app_state, user, target, xtream_output, &vod_pli) { + return value; + } } + Err(err) => error!("Failed to persist video info: {err}") } } } @@ -112,13 +121,15 @@ pub async fn get_xtream_stream_info(client: &reqwest::Client, let group = pli.get_group(); let series_name = pli.get_name(); - if let Ok(storage_path) = get_input_storage_path(&input.name, working_dir) { - if let Ok(info) = serde_json::from_str::(&content) { + match serde_json::from_str::(&content) { + Ok(info) => { // parse series info let series_stream_props = SeriesStreamProperties::from_info(&info, pli); - // update input db - if let Err(err) = persists_input_series_info(app_config, &storage_path, cluster, &input.name, provider_id, &series_stream_props).await { - error!("Failed to persist series info for input {}: {err}", &input.name); + if let Ok(storage_path) = get_input_storage_path(&input.name, working_dir) { + // update input db + if let Err(err) = persists_input_series_info(app_config, &storage_path, cluster, &input.name, provider_id, &series_stream_props).await { + error!("Failed to persist series info for input {}: {err}", &input.name); + } } if let Some(mut episodes) = parse_xtream_series_info(&pli.get_uuid(), &series_stream_props, &group, &series_name, input) { @@ -137,10 +148,11 @@ pub async fn get_xtream_stream_info(client: &reqwest::Client, for episode in &mut episodes { episode.header.virtual_id = target_id_mapping.get_and_update_virtual_id(&episode.header.uuid, provider_id, episode.header.item_type, parent_id); episode.header.category_id = category_id; - provider_series.entry(pli.parent_code.clone()) + let episode_provider_id = episode.header.get_provider_id().unwrap_or(0); + provider_series.entry(pli.get_uuid().to_string()) .or_default() .push(ProviderEpisodeKey { - provider_id: episode.header.get_provider_id().unwrap_or(0), + provider_id: episode_provider_id, virtual_id: episode.header.virtual_id, }); if target.use_memory_cache { @@ -156,38 +168,67 @@ pub async fn get_xtream_stream_info(client: &reqwest::Client, } } } + if let Err(err) = target_id_mapping.persist() { + error!("Failed to persist target id mapping: {err}"); + } } - if !provider_series.is_empty() { - let mut series_pli = pli.clone(); - series_pli.additional_properties = Some(StreamProperties::Series(Box::new(series_stream_props))); - rewrite_provider_series_info_episode_virtual_id(&mut series_pli, &provider_series); - if let Err(err) = write_playlist_item_to_file(app_config, &target.name, &series_pli).await { - error!("Failed to persist series stream: {err}"); - } - app_state.playlists.update_playlist_items(target, vec![&series_pli]).await; + let xtream_episodes: Vec = episodes.iter().map(XtreamPlaylistItem::from).collect(); + if let Err(err) = write_playlist_batch_item_upsert( + app_config, + &target.name, + XtreamCluster::Series, + &xtream_episodes).await { + error!("Failed to persist playlist batch item update: {err}"); } if target.use_memory_cache && !in_memory_updates.is_empty() { app_state.playlists.insert_playlist_items(target, episodes).await; app_state.playlists.update_target_id_mapping(target, in_memory_updates).await; } + + if !provider_series.is_empty() { + let mut series_pli = pli.clone(); + series_pli.additional_properties = Some(StreamProperties::Series(Box::new(series_stream_props))); + rewrite_provider_series_info_episode_virtual_id(&mut series_pli, &provider_series); + if let Err(err) = write_playlist_item_update(app_config, &target.name, &series_pli).await { + error!("Failed to persist series stream: {err}"); + } + app_state.playlists.update_playlist_items(target, vec![&series_pli]).await; + + if let Some(value) = xtream_resolve_stream_info(app_state, user, target, xtream_output, &series_pli) { + return value; + } + } } } } } + Err(err) => { + error!("Failed to persist series info: {err}"); + } } } } } - - return Ok(content); } Err(str_to_io_error(&format!("Cant find stream with id: {}/{}/{}", target.name.replace(' ', "_").as_str(), &cluster, pli.get_virtual_id()))) } +fn xtream_resolve_stream_info(app_state: &Arc, user: &ProxyUserCredentials, + target: &ConfigTarget, xtream_output: &XtreamTargetOutput, + pli: &XtreamPlaylistItem) -> Option> { + let app_config = &app_state.app_config; + let server_info = app_config.get_user_server_info(user); + let options = xtream_mapping_option_from_target_options(target, xtream_output, app_config, user, Some(server_info.get_base_url().as_str())); + if let Some(content) = pli.get_resolved_info_document(&options) { + return Some(serde_json::to_string(&content).map_err(to_io_error)); + } + None +} + fn get_skip_cluster(input: &ConfigInput) -> Vec { let mut skip_cluster = vec![]; if let Some(input_options) = &input.options { diff --git a/shared/Cargo.toml b/shared/Cargo.toml index 4fadefcdc..74792d541 100644 --- a/shared/Cargo.toml +++ b/shared/Cargo.toml @@ -25,6 +25,7 @@ bytes = "1" ciborium = "0.2.2" hex = "0.4.3" uuid = { version = "1.19.0", features = ["v4"] } +lz4_flex = "0.12.0" [target.'cfg(target_arch = "wasm32")'.dependencies] js-sys = "0.3.82" diff --git a/shared/src/model/stream_properties.rs b/shared/src/model/stream_properties.rs index fe1691535..e107dfb1d 100644 --- a/shared/src/model/stream_properties.rs +++ b/shared/src/model/stream_properties.rs @@ -1,7 +1,7 @@ use crate::model::info_doc_utils::InfoDocUtils; use crate::model::{PlaylistEntry, PlaylistItemType, VirtualId, XtreamCluster, XtreamMappingOptions, XtreamSeriesInfo, XtreamVideoInfo}; use crate::utils::{deserialize_as_option_string, deserialize_as_string, - deserialize_as_string_array, deserialize_json_as_opt_string, + deserialize_as_string_array, deserialize_json_as_opt_string, serialize_json_as_opt_string, deserialize_number_from_string, deserialize_number_from_string_or_zero, string_default_on_null, string_or_number_u32}; use serde::{Deserialize, Serialize}; @@ -62,9 +62,9 @@ pub struct VideoStreamDetailProperties { pub backdrop_path: Option>, pub duration_secs: Option, pub duration: Option, - #[serde(default, deserialize_with = "deserialize_json_as_opt_string")] + #[serde(default, serialize_with = "serialize_json_as_opt_string", deserialize_with = "deserialize_json_as_opt_string")] pub video: Option, - #[serde(default, deserialize_with = "deserialize_json_as_opt_string")] + #[serde(default, serialize_with = "serialize_json_as_opt_string", deserialize_with = "deserialize_json_as_opt_string")] pub audio: Option, #[serde(default)] pub bitrate: u32, @@ -140,9 +140,9 @@ pub struct SeriesStreamDetailEpisodeProperties { pub bitrate: u32, #[serde(default, deserialize_with = "deserialize_number_from_string")] pub rating: Option, - #[serde(default, deserialize_with = "deserialize_json_as_opt_string")] + #[serde(default, serialize_with = "serialize_json_as_opt_string", deserialize_with = "deserialize_json_as_opt_string")] pub video: Option, - #[serde(default, deserialize_with = "deserialize_json_as_opt_string")] + #[serde(default, serialize_with = "serialize_json_as_opt_string", deserialize_with = "deserialize_json_as_opt_string")] pub audio: Option, } @@ -205,9 +205,9 @@ pub struct EpisodeStreamProperties { pub movie_image: String, #[serde(default, deserialize_with = "string_default_on_null")] pub container_extension: String, - #[serde(default, deserialize_with = "deserialize_json_as_opt_string")] + #[serde(default, serialize_with = "serialize_json_as_opt_string", deserialize_with = "deserialize_json_as_opt_string")] pub video: Option, - #[serde(default, deserialize_with = "deserialize_json_as_opt_string")] + #[serde(default, serialize_with = "serialize_json_as_opt_string", deserialize_with = "deserialize_json_as_opt_string")] pub audio: Option, } diff --git a/shared/src/model/xtream.rs b/shared/src/model/xtream.rs index b221cfb37..494117167 100644 --- a/shared/src/model/xtream.rs +++ b/shared/src/model/xtream.rs @@ -1,7 +1,7 @@ use crate::utils::{deserialize_as_option_string, deserialize_as_string_array, - deserialize_json_as_opt_string, deserialize_number_from_string, - deserialize_number_from_string_or_zero, - string_default_on_null, string_or_number_u32}; + deserialize_number_from_string, deserialize_number_from_string_or_zero, + deserialize_as_string, string_default_on_null, string_or_number_u32, + deserialize_json_as_opt_string, serialize_json_as_opt_string}; use serde::ser::SerializeMap; use serde::{Deserialize, Deserializer, Serialize, Serializer}; use serde_json::Value; @@ -18,7 +18,7 @@ pub struct XtreamVideoInfoMovieData { pub direct_source: String, #[serde(default, deserialize_with = "deserialize_as_option_string")] pub custom_sid: Option, - #[serde(default)] + #[serde(default, deserialize_with = "deserialize_as_string")] pub added: String, #[serde(default)] pub container_extension: String, @@ -26,14 +26,16 @@ pub struct XtreamVideoInfoMovieData { #[derive(Serialize, Deserialize, Debug, Clone, PartialEq)] pub struct XtreamVideoInfoInfo { - pub kinopoisk_url: Option, #[serde(default)] + pub kinopoisk_url: Option, + #[serde(default, deserialize_with = "deserialize_as_string")] pub tmdb_id: String, // is in get_vod_streams #[serde(default)] pub name: String, // is in get_vod_streams pub o_name: Option, pub cover_big: Option, pub movie_image: Option, + #[serde(default, deserialize_with = "deserialize_as_option_string")] pub releasedate: Option, #[serde(default, deserialize_with = "deserialize_number_from_string")] pub episode_run_time: Option, @@ -43,7 +45,9 @@ pub struct XtreamVideoInfoInfo { pub cast: Option, pub description: Option, pub plot: Option, + #[serde(default, deserialize_with = "deserialize_as_option_string")] pub age: Option, + #[serde(default, deserialize_with = "deserialize_as_option_string")] pub mpaa_rating: Option, #[serde(default)] pub rating_count_kinopoisk: u32, @@ -51,15 +55,19 @@ pub struct XtreamVideoInfoInfo { pub genre: Option, #[serde(default, deserialize_with = "deserialize_as_string_array")] pub backdrop_path: Option>, + #[serde(default, deserialize_with = "deserialize_as_option_string")] pub duration_secs: Option, + #[serde(default, deserialize_with = "deserialize_as_option_string")] pub duration: Option, - #[serde(default, deserialize_with = "deserialize_json_as_opt_string")] + #[serde(default, serialize_with = "serialize_json_as_opt_string", deserialize_with = "deserialize_json_as_opt_string")] pub video: Option, - #[serde(default, deserialize_with = "deserialize_json_as_opt_string")] + #[serde(default, serialize_with = "serialize_json_as_opt_string", deserialize_with = "deserialize_json_as_opt_string")] pub audio: Option, #[serde(default)] pub bitrate: u32, + #[serde(default, deserialize_with = "deserialize_as_option_string")] pub runtime: Option, + #[serde(default, deserialize_with = "deserialize_as_option_string")] pub status: Option, } @@ -107,7 +115,7 @@ pub struct XtreamSeriesInfoInfo { pub director: String, #[serde(default, deserialize_with = "string_default_on_null")] pub genre: String, - #[serde(default, alias = "releaseDate", alias = "releasedate", deserialize_with = "string_default_on_null")] + #[serde(default, deserialize_with = "string_default_on_null")] pub release_date: String, #[serde(default, deserialize_with = "string_default_on_null")] pub last_modified: String, diff --git a/shared/src/utils/serde_utils.rs b/shared/src/utils/serde_utils.rs index 4307296aa..3afb30df4 100644 --- a/shared/src/utils/serde_utils.rs +++ b/shared/src/utils/serde_utils.rs @@ -1,6 +1,9 @@ use crate::error::to_io_error; +use base64::engine::general_purpose; +use base64::Engine; use chrono::{NaiveDateTime, ParseError, TimeZone, Utc}; -use serde::{Deserialize, Deserializer}; +use serde::de::Error; +use serde::{Deserialize, Deserializer, Serializer}; use serde_json::Value; use std::io; @@ -260,11 +263,76 @@ pub fn parse_timestamp(value: &str) -> Result, ParseError> { let timestamp = Utc.from_utc_datetime(&dt).timestamp(); Ok(Some(timestamp)) } +// +// pub fn deserialize_json_as_opt_string<'de, D>(deserializer: D) -> Result, D::Error> +// where +// D: Deserializer<'de>, +// { +// let val: Value = Deserialize::deserialize(deserializer)?; +// Ok(Some(val.to_string())) +// } +// +// pub fn serialize_json_opt_string<'de, S>(value: &String, serializer: S) -> Result, S::Error> +// where +// S: Serializer, +// { +// let json: Value = serde_json::from_str(value).map_err(serde::ser::Error::custom)?; +// json.serialize(serializer) +// } -pub fn deserialize_json_as_opt_string<'de, D>(deserializer: D) -> Result, D::Error> +pub fn serialize_json_as_opt_string(value: &Option, + serializer: S, +) -> Result +where + S: Serializer, +{ + let s = match value { + Some(s) => s, + None => return serializer.serialize_none(), + }; + + let bytes = s.as_bytes(); + let compressed = lz4_flex::compress_prepend_size(bytes); + let encoded = general_purpose::STANDARD_NO_PAD.encode(compressed); + + serializer.serialize_some(&encoded) +} + +pub fn deserialize_json_as_opt_string<'de, D>( + deserializer: D, +) -> Result, D::Error> where D: Deserializer<'de>, { - let val: Value = Deserialize::deserialize(deserializer)?; - Ok(Some(val.to_string())) + let opt: Option = Option::deserialize(deserializer)?; + + match opt { + None => Ok(None), + Some(Value::Array(arr)) if arr.is_empty() => Ok(None), + Some(Value::Object(obj)) if obj.is_empty() => Ok(None), + Some(Value::String(s)) => { + if s.is_empty() { + Ok(None) + } else { + let compressed = match base64::engine::general_purpose::STANDARD_NO_PAD.decode(&s) { + Ok(bytes) => bytes, + Err(_) => return Ok(Some(s)), + }; + + let decompressed = match lz4_flex::decompress_size_prepended(&compressed) { + Ok(bytes) => bytes, + Err(_) => return Ok(None), + }; + + match String::from_utf8(decompressed) { + Ok(text) => Ok(Some(text)), + Err(_) => Ok(None), + } + } + } + Some(Value::Null) + | Some(Value::Number(_)) + | Some(Value::Bool(_)) => Ok(None), + Some(v) => Ok(Some(serde_json::to_string(&v).map_err(D::Error::custom)?)), + } }