diff --git a/Cargo.lock b/Cargo.lock index 5122f908d..f2f4e0b83 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -506,9 +506,9 @@ checksum = "1fd0f2584146f6f2ef48085050886acf353beff7305ebd1ae69500e27c67f64b" [[package]] name = "bytes" -version = "1.7.2" +version = "1.8.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "428d9aa8fbc0670b7b8d6030a7fadd0f86151cae55e4dbbece15f3780a3dfaf3" +checksum = "9ac0150caa2ae65ca5bd83f25c7de183dea78d4d366469f148435e2acfbad0da" [[package]] name = "bytestring" @@ -527,9 +527,9 @@ checksum = "a2698f953def977c68f935bb0dfa959375ad4638570e969e2f1e9f433cbf1af6" [[package]] name = "cc" -version = "1.1.30" +version = "1.1.31" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "b16803a61b81d9eabb7eae2588776c4c1e584b738ede45fdbb4c972cec1e9945" +checksum = "c2e7962b54006dcfcc61cb72735f4d89bb97061dd6a7ed882ec6b8ee53714c6f" dependencies = [ "jobserver", "libc", @@ -602,6 +602,15 @@ version = "1.0.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d3fd119d74b830634cea2a0f58bbd0d54540518a14397557951e79340abc28c0" +[[package]] +name = "compressed_string" +version = "1.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d3191fb9cd59eb0195c223d4f1d8894d2078957ee44f4940f059d11060a6d0e7" +dependencies = [ + "flate2", +] + [[package]] name = "concurrent-queue" version = "2.5.0" @@ -760,9 +769,9 @@ dependencies = [ [[package]] name = "encoding_rs" -version = "0.8.34" +version = "0.8.35" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "b45de904aa0b010bce2ab45264d0631681847fa7b6f2eaa7dab7619943bc4f59" +checksum = "75030f3c4f45dafd7586dd6780965a8c7e8e285a5ecb86713e63a79c5b2766f3" dependencies = [ "cfg-if", ] @@ -1259,9 +1268,9 @@ dependencies = [ [[package]] name = "impl-more" -version = "0.1.7" +version = "0.1.8" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e658178c10c747241199382079c0f195ce229866fbf4aa0d46fa6107fe33d2ec" +checksum = "aae21c3177a27788957044151cc2800043d127acaa460a47ebb9b84dfa2c6aa0" [[package]] name = "indexmap" @@ -1370,9 +1379,9 @@ checksum = "d4345964bb142484797b161f473a503a434de77149dd8c7427788c6e13379388" [[package]] name = "libc" -version = "0.2.159" +version = "0.2.161" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "561d97a539a36e26a9a5fad1ea11a3039a67714694aaa379433e580854bc3dc5" +checksum = "8e9489c2807c139ffd9c1794f4af0ebe86a828db53ecdc7fea2111d0fed085d1" [[package]] name = "libnghttp2-sys" @@ -1445,11 +1454,11 @@ dependencies = [ "actix-server", "actix-web", "actix-web-httpauth", - "base64 0.22.1", "bincode", "blake3", "chrono", "clap", + "compressed_string", "cron", "enum-iterator", "env_logger", @@ -1607,9 +1616,9 @@ checksum = "1261fe7e33c73b354eab43b1273a57c8f967d0391e80353e51f764ac02cf6775" [[package]] name = "openssl" -version = "0.10.67" +version = "0.10.68" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "7b8cefcf97f41316955f9294cd61f639bdcfa9f2f230faac6cb896aa8ab64704" +checksum = "6174bc48f102d208783c2c84bf931bb75927a617866870de8a4ea85597f871f5" dependencies = [ "bitflags 2.6.0", "cfg-if", @@ -1639,9 +1648,9 @@ checksum = "ff011a302c396a5197692431fc1948019154afc178baf7d8e37367442a4601cf" [[package]] name = "openssl-src" -version = "300.3.2+3.3.2" +version = "300.4.0+3.4.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "a211a18d945ef7e648cc6e0058f4c548ee46aab922ea203e0d30e966ea23647b" +checksum = "a709e02f2b4aca747929cca5ed248880847c650233cf8b8cdc48f40aaf4898a6" dependencies = [ "cc", ] @@ -1785,18 +1794,18 @@ dependencies = [ [[package]] name = "pin-project" -version = "1.1.6" +version = "1.1.7" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "baf123a161dde1e524adf36f90bc5d8d3462824a9c43553ad07a8183161189ec" +checksum = "be57f64e946e500c8ee36ef6331845d40a93055567ec57e8fae13efd33759b95" dependencies = [ "pin-project-internal", ] [[package]] name = "pin-project-internal" -version = "1.1.6" +version = "1.1.7" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "a4502d8515ca9f32f1fb543d987f63d95a14934883db45bdb48060b6b69257f8" +checksum = "3c0f5fad0874fc7abcd4d750e76917eaebbecaa2c20bde22e1dbeeba8beb758c" dependencies = [ "proc-macro2", "quote", @@ -1805,9 +1814,9 @@ dependencies = [ [[package]] name = "pin-project-lite" -version = "0.2.14" +version = "0.2.15" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "bda66fc9667c18cb2758a2ac84d1167245054bcf85d5d1aaa6923f45801bdd02" +checksum = "915a1e146535de9163f3987b8944ed8cf49a18bb0056bcebcdcece385cece4ff" [[package]] name = "pin-utils" @@ -1854,9 +1863,9 @@ dependencies = [ [[package]] name = "proc-macro2" -version = "1.0.87" +version = "1.0.89" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "b3e4daa0dcf6feba26f985457cdf104d4b4256fc5a09547140f3631bb076b19a" +checksum = "f139b0662de085916d1fb67d2b4169d1addddda1919e696f3252b740b629986e" dependencies = [ "unicode-ident", ] @@ -2149,9 +2158,9 @@ dependencies = [ [[package]] name = "rustls" -version = "0.23.14" +version = "0.23.15" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "415d9944693cb90382053259f89fbb077ea730ad7273047ec63b19bc9b160ba8" +checksum = "5fbb44d7acc4e873d613422379f69f237a1b141928c02f6bc6ccfddddc2d7993" dependencies = [ "once_cell", "ring", @@ -2239,18 +2248,18 @@ checksum = "61697e0a1c7e512e84a621326239844a24d8207b4669b41bc18b32ea5cbf988b" [[package]] name = "serde" -version = "1.0.210" +version = "1.0.213" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "c8e3592472072e6e22e0a54d5904d9febf8508f65fb8552499a1abc7d1078c3a" +checksum = "3ea7893ff5e2466df8d720bb615088341b295f849602c6956047f8f80f0e9bc1" dependencies = [ "serde_derive", ] [[package]] name = "serde_derive" -version = "1.0.210" +version = "1.0.213" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "243902eda00fad750862fc144cea25caca5e20d615af0a81bee94ca738f1df1f" +checksum = "7e85ad2009c50b58e87caa8cd6dac16bdf511bbfb7af6c33df902396aa480fa5" dependencies = [ "proc-macro2", "quote", @@ -2259,9 +2268,9 @@ dependencies = [ [[package]] name = "serde_json" -version = "1.0.128" +version = "1.0.132" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6ff5456707a1de34e7e37f2a6fd3d3f808c318259cbd01ab6377795054b483d8" +checksum = "d726bfaff4b320266d395898905d0eba0345aae23b54aee3a737e260fd46db03" dependencies = [ "itoa", "memchr", @@ -2405,9 +2414,9 @@ checksum = "13c2bddecc57b384dee18652358fb23172facb8a2c51ccc10d74c157bdea3292" [[package]] name = "syn" -version = "2.0.79" +version = "2.0.85" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "89132cd0bf050864e1d38dc3bbc07a0eb8e7530af26344d3d2bbbef83499f590" +checksum = "5023162dfcd14ef8f32034d8bcd4cc5ddc61ef7a247c024a33e24e1f24d21b56" dependencies = [ "proc-macro2", "quote", @@ -2459,18 +2468,18 @@ dependencies = [ [[package]] name = "thiserror" -version = "1.0.64" +version = "1.0.65" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "d50af8abc119fb8bb6dbabcfa89656f46f84aa0ac7688088608076ad2b459a84" +checksum = "5d11abd9594d9b38965ef50805c5e469ca9cc6f197f883f717e0269a3057b3d5" dependencies = [ "thiserror-impl", ] [[package]] name = "thiserror-impl" -version = "1.0.64" +version = "1.0.65" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "08904e7672f5eb876eaaf87e0ce17857500934f4981c4a0ab2b4aa98baac7fc3" +checksum = "ae71770322cbd277e69d762a16c444af02aa0575ac0d174f0b9562d3b37f8602" dependencies = [ "proc-macro2", "quote", @@ -2525,9 +2534,9 @@ checksum = "1f3ccbac311fea05f86f61904b462b55fb3df8837a366dfc601a0161d0532f20" [[package]] name = "tokio" -version = "1.40.0" +version = "1.41.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e2b070231665d27ad9ec9b8df639893f46727666c6767db40317fbe920a5d998" +checksum = "145f3413504347a2be84393cc8a7d2fb4d863b375909ea59f2158261aa258bbb" dependencies = [ "backtrace", "bytes", @@ -2642,12 +2651,9 @@ checksum = "2896d95c02a80c6d6a5d6e953d479f5ddf2dfdb6a244441010e373ac0fb88971" [[package]] name = "unicase" -version = "2.7.0" +version = "2.8.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "f7d2d4dafb69621809a81864c9c1b864479e1235c0dd4e199924b9742439ed89" -dependencies = [ - "version_check", -] +checksum = "7e51b68083f157f853b6379db119d1c1be0e6e4dec98101079dec41f6f5cf6df" [[package]] name = "unicode-bidi" diff --git a/Cargo.toml b/Cargo.toml index ba849560a..2d1b392fc 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -14,10 +14,10 @@ strip = true [dependencies] serde = { version = "1.0", features = ["derive", "rc"] } -serde_yaml = "0.9.33" +serde_yaml = "0.9.34" serde_json = "1" quick-xml = { version = "0.36", features = ["serialize"] } -regex = "1.10" +regex = "1.11" clap = { version = "4", features = ["derive"] } url = "2.5" reqwest = { version = "0", features = ["blocking", "json", "stream", "rustls-tls"] } @@ -49,5 +49,5 @@ rand = "0.8" rpassword = "7.3" flate2 = "1" time = "0.3" -blake3 = "1.5.4" -base64 = "0.22.1" +blake3 = "1.5" +compressed_string = "1.0.0" diff --git a/src/api/m3u_api.rs b/src/api/m3u_api.rs index 94d72fd98..1c2bf70c9 100644 --- a/src/api/m3u_api.rs +++ b/src/api/m3u_api.rs @@ -5,6 +5,7 @@ use crate::api::api_utils::{get_user_target, get_user_target_by_credentials, str use crate::api::api_model::{AppState, UserApiRequest}; use crate::model::config::TargetType; use crate::repository::m3u_repository::{m3u_get_file_paths, m3u_get_item_for_stream_id, m3u_load_rewrite_playlist}; +use crate::repository::storage::get_target_storage_path; async fn m3u_api( api_req: web::Query, @@ -38,15 +39,23 @@ async fn m3u_api_stream( if let Ok(m3u_stream_id) = stream_id.parse::() { if let Some((_user, target)) = get_user_target_by_credentials(&username, &password, &api_req, &app_state) { if target.has_output(&TargetType::M3u) { - if let Some((m3u_path, idx_path)) = m3u_get_file_paths(&app_state.config, target) { - match m3u_get_item_for_stream_id(m3u_stream_id, &m3u_path, &idx_path) { - Ok(m3u_item) => { - return stream_response(m3u_item.url.as_str(), &req, None).await; - } - Err(err) => { - error!("Failed to get m3u url: {}", err); + + 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) { + Ok(m3u_item) => { + return stream_response(m3u_item.url.as_str(), &req, None).await; + } + Err(err) => { + error!("Failed to get m3u url: {}", err); + } + } } } + None => { + error!("Failed to get target path for {}", target.name); + } } } } diff --git a/src/api/xmltv_api.rs b/src/api/xmltv_api.rs index 4a54e9427..ad310cff9 100644 --- a/src/api/xmltv_api.rs +++ b/src/api/xmltv_api.rs @@ -17,6 +17,7 @@ use crate::model::api_proxy::{ProxyType, ProxyUserCredentials}; use crate::model::config::{Config, ConfigInput, ConfigTarget}; use crate::model::config::TargetType; use crate::repository::m3u_repository::m3u_get_epg_file_path; +use crate::repository::storage::get_target_storage_path; use crate::repository::xtream_repository::{xtream_get_epg_file_path, xtream_get_storage_path}; use crate::utils::{file_utils, request_utils}; @@ -51,10 +52,13 @@ fn get_epg_path_for_target_of_type(target_name: &str, file_path: Option } fn get_epg_path_for_target(config: &Config, target: &ConfigTarget) -> Option { + // TODO if we share the same virtual_id for epg, can we store an epg file for the target ? for output in &target.output { match output.target { TargetType::M3u => { - return get_epg_path_for_target_of_type(&target.name, m3u_get_epg_file_path(config, target)); + if let Some(target_path) = get_target_storage_path(config, &target.name) { + return get_epg_path_for_target_of_type(&target.name, m3u_get_epg_file_path(&target_path)); + } } TargetType::Xtream => { if let Some(storage_path) = xtream_get_storage_path(config, &target.name) { diff --git a/src/api/xtream_api.rs b/src/api/xtream_api.rs index 972701d69..12fb38130 100644 --- a/src/api/xtream_api.rs +++ b/src/api/xtream_api.rs @@ -263,7 +263,7 @@ async fn xtream_player_api_streaming_timeshift( fn get_xtream_vod_info(target: &ConfigTarget, pli: &XtreamPlaylistItem, content: &str) -> Result { if let Ok(mut doc) = serde_json::from_str::>(content) { if let Some(Value::Object(movie_data)) = doc.get_mut("movie_data") { - let stream_id = pli.stream_id; + let stream_id = pli.virtual_id; let category_id = pli.category_id; movie_data.insert("stream_id".to_string(), Value::Number(serde_json::value::Number::from(stream_id))); movie_data.insert("category_id".to_string(), Value::Number(serde_json::value::Number::from(category_id))); @@ -278,7 +278,7 @@ fn get_xtream_vod_info(target: &ConfigTarget, pli: &XtreamPlaylistItem, content: } } } - Err(Error::new(ErrorKind::Other, format!("Failed to get vod info for id {}", pli.stream_id))) + Err(Error::new(ErrorKind::Other, format!("Failed to get vod info for id {}", pli.virtual_id))) } fn get_xtream_series_info(config: &Config, target: &ConfigTarget, pli: &XtreamPlaylistItem, content: &str) -> Result { @@ -309,13 +309,13 @@ fn get_xtream_series_info(config: &Config, target: &ConfigTarget, pli: &XtreamPl } } if let Ok(result) = serde_json::to_string(&doc) { - let _ = xtream_repository::xtream_write_series_info(config, target.name.replace(' ', "_").as_str(), pli.stream_id, &new_id_to_provider_id_mapping, &result); + let _ = xtream_repository::xtream_write_series_info(config, target.name.replace(' ', "_").as_str(), pli.virtual_id, &new_id_to_provider_id_mapping, &result); return Ok(result); } } } } - Err(Error::new(ErrorKind::Other, format!("Failed to get series info for id {}", pli.stream_id))) + Err(Error::new(ErrorKind::Other, format!("Failed to get series info for id {}", pli.virtual_id))) } async fn xtream_get_stream_info_content(info_url: &str, input: &ConfigInput) -> Result { @@ -325,7 +325,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 { - if let Ok(content) = xtream_repository::xtream_load_series_info(config, target.name.replace(' ', "_").as_str(), pli.stream_id) { + if let Ok(content) = xtream_repository::xtream_load_series_info(config, target.name.replace(' ', "_").as_str(), pli.virtual_id) { return Ok(content); } } @@ -343,7 +343,7 @@ async fn xtream_get_stream_info(config: &Config, input: &ConfigInput, target: &C } Err(Error::new(std::io::ErrorKind::Other, format!("Cant find stream with id: {}/{}/{}", - target.name.replace(' ', "_").as_str(), &cluster, pli.stream_id))) + target.name.replace(' ', "_").as_str(), &cluster, pli.virtual_id))) } async fn xtream_get_stream_info_response(app_state: &AppState, user: &ProxyUserCredentials, @@ -549,13 +549,8 @@ async fn xtream_player_api( } } _ => { - if api_req.action.is_empty() { - debug!("Paremeter action is empty!"); - HttpResponse::Unauthorized().finish() - } else { - debug!("cant find user!"); - HttpResponse::BadRequest().finish() - } + debug!("{}", if api_req.action.is_empty() { "Paremeter action is empty!" } else { "cant find user!" }); + HttpResponse::BadRequest().finish() } } } diff --git a/src/main.rs b/src/main.rs index bc28432bd..ca68cf275 100644 --- a/src/main.rs +++ b/src/main.rs @@ -147,11 +147,10 @@ fn init_logger(log_level: &str) { let mut log_builder = Builder::from_default_env(); if log_level.contains('=') { - let pairs: Vec<&str> = log_level.split(',').collect(); - for pair in pairs { - let kv: Vec<&str> = pair.split('=').collect(); - if kv.len() == 2 { - log_builder.filter_module(kv[0].trim(), get_log_level(kv[1].trim())); + for pair in log_level.split(',').filter(|s| s.contains('=')) { + let mut kv_iter = pair.split('=').map(str::trim); + if let (Some(module), Some(level)) = (kv_iter.next(), kv_iter.next()) { + log_builder.filter_module(module, get_log_level(level)); } } } else { diff --git a/src/model/playlist.rs b/src/model/playlist.rs index 7ab5c4ebb..ec883a60f 100644 --- a/src/model/playlist.rs +++ b/src/model/playlist.rs @@ -2,9 +2,6 @@ use std::cell::RefCell; use std::cmp::PartialEq; use std::fmt::{Display, Formatter}; use std::rc::Rc; -use base64::Engine; -use blake3::Hasher; -use base64::engine::general_purpose; use serde::{Deserialize, Serialize}; use serde_json::Value; use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; @@ -12,8 +9,8 @@ use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; use crate::model::config::{ConfigInput, ConfigTarget}; use crate::model::xmltv::TVGuide; use crate::model::xtream::{xtream_playlistitem_to_document, XtreamMappingOptions}; -use crate::utils::default_utils::{default_as_false, default_as_zero_u16, default_as_zero_u32, default_playlist_item_type, default_stream_cluster}; -use crate::utils::request_utils::get_base_url; +use crate::repository::storage::hash_string; +use crate::utils::default_utils::{default_as_false, default_as_zero_u16, default_as_zero_u32}; // https://de.wikipedia.org/wiki/M3U // https://siptv.eu/howto/playlist.html @@ -44,6 +41,12 @@ pub(crate) enum XtreamCluster { Series = 3, } +impl Default for XtreamCluster { + fn default() -> Self { + XtreamCluster::Live + } +} + impl Display for XtreamCluster { fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result { write!(f, "{}", match self { @@ -62,6 +65,22 @@ pub(crate) enum PlaylistItemType { SeriesInfo = 4, } +impl Default for PlaylistItemType { + fn default() -> Self { + PlaylistItemType::Live + } +} + +impl From for PlaylistItemType { + fn from(xtream_cluster: XtreamCluster) -> Self { + match xtream_cluster { + XtreamCluster::Live => PlaylistItemType::Live, + XtreamCluster::Video => PlaylistItemType::Movie, + XtreamCluster::Series => PlaylistItemType::SeriesInfo, + } + } +} + impl Display for PlaylistItemType { fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result { write!(f, "{}", match self { @@ -78,11 +97,11 @@ pub(crate) trait FieldAccessor { fn set_field(&mut self, field: &str, value: &str) -> bool; } -#[derive(Debug, Clone, Serialize, Deserialize)] +#[derive(Debug, Clone, Serialize, Deserialize, Default)] pub(crate) struct PlaylistItemHeader { - pub uuid: Rc, // calculated - pub stream_id: Rc, // virtual id + pub uuid: Rc<[u8;32]>, // calculated pub id: Rc, // provider id + pub virtual_id: u32, // virtual id pub name: Rc, pub chno: Rc, pub logo: Rc, @@ -95,10 +114,9 @@ pub(crate) struct PlaylistItemHeader { pub rec: Rc, pub url: Rc, pub epg_channel_id: Option>, - #[serde(default = "default_stream_cluster")] pub xtream_cluster: XtreamCluster, pub additional_properties: Option, - #[serde(default = "default_playlist_item_type", skip_serializing, skip_deserializing)] + #[serde(skip_serializing, skip_deserializing)] pub item_type: PlaylistItemType, #[serde(default = "default_as_false", skip_serializing, skip_deserializing)] pub series_fetched: bool, // only used for series_info @@ -111,20 +129,10 @@ pub(crate) struct PlaylistItemHeader { impl PlaylistItemHeader { pub(crate) fn gen_uuid(&mut self) { - //let base_url = get_base_url(&self.url).unwrap_or_else(|| self.url.to_string()); - // Create a Blake3 hasher - let mut hasher = Hasher::new(); - // the url should be different for each entry. - hasher.update(self.url.as_bytes()); - // Finalize and get the hash result - let hash_bytes = hasher.finalize(); - // Encode the reduced 128-bit hash to Base62 - let mut encoded_key = String::new(); - general_purpose::STANDARD.encode_string(&hash_bytes.as_bytes(), &mut encoded_key); - self.uuid = Rc::new(encoded_key); + self.uuid = Rc::new(hash_string(&self.url)) } - pub(crate) fn get_uuid(&self) -> &str { - self.uuid.as_ref() + pub(crate) fn get_uuid(&self) -> &Rc<[u8; 32]> { + &self.uuid } } @@ -172,11 +180,11 @@ macro_rules! generate_field_accessor_impl_for_playlist_item_header { } } -generate_field_accessor_impl_for_playlist_item_header!(id, stream_id, name, chno, logo, logo_small, group, title, parent_code, audio_track, time_shift, rec, url;); +generate_field_accessor_impl_for_playlist_item_header!(id, /*virtual_id,*/ name, chno, logo, logo_small, group, title, parent_code, audio_track, time_shift, rec, url;); #[derive(Debug, Clone, Serialize, Deserialize)] pub(crate) struct M3uPlaylistItem { - pub stream_id: Rc, + pub virtual_id: u32, pub provider_id: Rc, pub name: Rc, pub chno: Rc, @@ -218,7 +226,7 @@ impl M3uPlaylistItem { #[derive(Debug, Clone, Serialize, Deserialize)] pub(crate) struct XtreamPlaylistItem { - pub stream_id: u32, + pub virtual_id: u32, pub provider_id: u32, pub name: Rc, pub logo: Rc, @@ -252,7 +260,7 @@ impl PlaylistItem { pub fn to_m3u(&self) -> M3uPlaylistItem { let header = self.header.borrow(); M3uPlaylistItem { - stream_id: Rc::clone(&header.stream_id), + virtual_id: header.virtual_id, provider_id: Rc::clone(&header.id), name: Rc::clone(&header.name), chno: Rc::clone(&header.chno), @@ -275,7 +283,7 @@ impl PlaylistItem { match header.id.parse::() { Ok(provider_id) => { Ok(XtreamPlaylistItem { - stream_id: header.stream_id.parse::().unwrap_or(0), + virtual_id: header.virtual_id, provider_id, name: Rc::clone(&header.name), logo: Rc::clone(&header.logo), @@ -312,7 +320,7 @@ pub(crate) struct PlaylistGroup { pub id: u32, pub title: Rc, pub channels: Vec, - #[serde(default = "default_stream_cluster", skip_serializing, skip_deserializing)] + #[serde(skip_serializing, skip_deserializing)] pub xtream_cluster: XtreamCluster, } @@ -321,5 +329,4 @@ impl PlaylistGroup { pub(crate) fn on_load(&mut self) { self.channels.iter().for_each(|pl| pl.header.borrow_mut().gen_uuid()); } -} - +} \ No newline at end of file diff --git a/src/model/xtream.rs b/src/model/xtream.rs index ba6522615..c1c199f3b 100644 --- a/src/model/xtream.rs +++ b/src/model/xtream.rs @@ -8,7 +8,7 @@ use serde_json::{Map, Value}; use crate::model::config::ConfigTargetOptions; use crate::model::playlist::{PlaylistItem, XtreamCluster, XtreamPlaylistItem}; -use crate::utils::default_utils::{default_as_empty_rc_str, default_as_empty_list}; +use crate::utils::default_utils::{default_as_empty_list, default_as_empty_rc_str}; const LIVE_STREAM_FIELDS: &[&str] = &[]; @@ -28,8 +28,8 @@ const SERIES_STREAM_FIELDS: &[&str] = &[ fn deserialize_number_from_string<'de, D, T: DeserializeOwned>( deserializer: D, ) -> Result, D::Error> - where - D: Deserializer<'de>, +where + D: Deserializer<'de>, { // we define a local enum type inside of the function // because it is untagged, serde will deserialize as the first variant @@ -73,8 +73,8 @@ fn value_to_string(v: &Value) -> Option { } fn deserialize_as_option_rc_string<'de, D>(deserializer: D) -> Result>, D::Error> - where - D: Deserializer<'de>, +where + D: Deserializer<'de>, { let value: Value = Deserialize::deserialize(deserializer)?; @@ -86,8 +86,8 @@ fn deserialize_as_option_rc_string<'de, D>(deserializer: D) -> Result(deserializer: D) -> Result, D::Error> - where - D: Deserializer<'de>, +where + D: Deserializer<'de>, { let value: Value = Deserialize::deserialize(deserializer)?; @@ -98,8 +98,8 @@ fn deserialize_as_rc_string<'de, D>(deserializer: D) -> Result, D::Er } fn deserialize_as_string_array<'de, D>(deserializer: D) -> Result>, D::Error> - where - D: Deserializer<'de>, +where + D: Deserializer<'de>, { Value::deserialize(deserializer).map(|v| match v { Value::Array(value) => Some(value_to_string_array(&value)), @@ -131,9 +131,9 @@ pub(crate) struct XtreamStream { #[serde(default, deserialize_with = "deserialize_as_rc_string")] pub category_id: Rc, #[serde(default, deserialize_with = "deserialize_number_from_string")] - pub stream_id: Option, + pub stream_id: Option, #[serde(default, deserialize_with = "deserialize_number_from_string")] - pub series_id: Option, + pub series_id: Option, #[serde(default = "default_as_empty_rc_str", deserialize_with = "deserialize_as_rc_string")] pub stream_icon: Rc, #[serde(default = "default_as_empty_rc_str", deserialize_with = "deserialize_as_rc_string")] @@ -217,10 +217,9 @@ macro_rules! add_i64_property_if_exists { } } - impl XtreamStream { - pub(crate) fn get_stream_id(&self) -> String { - self.stream_id.map_or_else(|| self.series_id.map_or_else(String::new, |seid| format!("{seid}")), |sid| format!("{sid}")) + pub(crate) fn get_stream_id(&self) -> u32 { + self.stream_id.unwrap_or(self.series_id.unwrap_or(0)) } pub(crate) fn get_additional_properties(&self) -> Option { @@ -315,6 +314,18 @@ pub(crate) struct XtreamSeriesInfoEpisode { pub direct_source: String, } +// impl XtreamSeriesInfoEpisode { +// pub(crate) fn get_id(&self) -> u32 { +// match self.id.parse::() { +// Ok(id) => id, +// Err(_) => { +// error!("Failed to convert id to number {}", self.id); +// 0 +// } +// } +// } +// } + #[derive(Debug, Clone, Serialize, Deserialize)] pub(crate) struct XtreamSeriesInfo { pub seasons: Vec, @@ -358,7 +369,7 @@ impl XtreamMappingOptions { pub fn from_target_options(options: Option<&ConfigTargetOptions>) -> Self { let (skip_live_direct_source, skip_video_direct_source, skip_series_direct_source) = options .map_or((false, false, false), |o| (o.xtream_skip_live_direct_source, - o.xtream_skip_video_direct_source, o.xtream_skip_series_direct_source)); + o.xtream_skip_video_direct_source, o.xtream_skip_series_direct_source)); Self { skip_live_direct_source, skip_video_direct_source, @@ -414,7 +425,7 @@ fn append_prepared_series_properties(add_props: Option<&Map>, doc } pub(crate) fn xtream_playlistitem_to_document(pli: &XtreamPlaylistItem, options: &XtreamMappingOptions) -> serde_json::Value { - let stream_id_value = Value::Number(serde_json::Number::from(pli.stream_id)); + let stream_id_value = Value::Number(serde_json::Number::from(pli.virtual_id)); let mut document = serde_json::Map::from_iter([ ("category_id".to_string(), Value::String(format!("{}", &pli.category_id))), ("category_ids".to_string(), Value::Array(Vec::from([Value::Number(serde_json::Number::from(pli.category_id))]))), diff --git a/src/processing/m3u_parser.rs b/src/processing/m3u_parser.rs index 6c4b05231..ba1cb96e0 100644 --- a/src/processing/m3u_parser.rs +++ b/src/processing/m3u_parser.rs @@ -4,7 +4,6 @@ use std::rc::Rc; use crate::model::config::{Config, ConfigInput}; use crate::model::playlist::{PlaylistGroup, PlaylistItem, PlaylistItemHeader, PlaylistItemType, XtreamCluster}; -use crate::utils::default_utils::{default_as_empty_rc_str, default_playlist_item_type, default_stream_cluster}; use crate::utils::string_utils; #[inline] @@ -76,27 +75,10 @@ fn skip_digit(it: &mut std::str::Chars) -> Option { fn create_empty_playlistitem_header(input_id: u16, url: &str) -> PlaylistItemHeader { PlaylistItemHeader { - uuid: default_as_empty_rc_str(), - id: default_as_empty_rc_str(), - stream_id: default_as_empty_rc_str(), - name: default_as_empty_rc_str(), - chno: default_as_empty_rc_str(), - logo: default_as_empty_rc_str(), - logo_small: default_as_empty_rc_str(), - group: default_as_empty_rc_str(), - title: default_as_empty_rc_str(), - parent_code: default_as_empty_rc_str(), - audio_track: default_as_empty_rc_str(), - time_shift: default_as_empty_rc_str(), - rec: default_as_empty_rc_str(), url: Rc::new(url.to_owned()), - epg_channel_id: None, - item_type: default_playlist_item_type(), - xtream_cluster: default_stream_cluster(), - additional_properties: None, - series_fetched: false, category_id: 0, input_id, + ..Default::default() } } @@ -147,7 +129,7 @@ fn process_header(input: &ConfigInput, video_suffixes: &Vec<&str>, content: &str plih.id = Rc::new(chanid); } } - plih.stream_id = Rc::clone(&plih.id); + // plih.virtual_id = plih.id; plih.epg_channel_id = Some(Rc::clone(&plih.id)); } @@ -174,7 +156,7 @@ fn process_header(input: &ConfigInput, video_suffixes: &Vec<&str>, content: &str pub(crate) fn extract_id_from_url(url: &str) -> Option { if let Some(filename) = url.split('/').last() { - return if let Some(index) = filename.rfind('.') { + return if let Some(index) = filename.rfind('.') { Some(filename[..index].to_string()) } else { Some(filename.to_string()) @@ -183,14 +165,17 @@ pub(crate) fn extract_id_from_url(url: &str) -> Option { None } -pub(crate) fn consume_m3u(cfg: &Config, input: &ConfigInput, lines: impl Iterator, mut visit: F) { +pub(crate) fn consume_m3u<'a, I, F: FnMut(PlaylistItem)>(cfg: &Config, input: &ConfigInput, lines: I, mut visit: F) +where + I: Iterator, +{ let mut header: Option = None; let mut group: Option = None; let video_suffixes = cfg.video.as_ref().unwrap().extensions.iter().map(String::as_str).collect::>(); for line in lines { if line.starts_with("#EXTINF") { - header = Some(line); + header = Some(String::from(line)); continue; } if line.starts_with("#EXTGRP") { @@ -219,29 +204,33 @@ pub(crate) fn consume_m3u(cfg: &Config, input: &ConfigIn } } -pub(crate) fn parse_m3u(cfg: &Config, input: &ConfigInput, lines: &[String]) -> Vec { +pub(crate) fn parse_m3u<'a, I>(cfg: &Config, input: &ConfigInput, lines: I) -> Vec +where + I: Iterator, +{ let mut sort_order: Vec> = vec![]; - let mut idx: usize = 0; + let mut sort_order_idx: usize = 0; let mut group_map: std::collections::HashMap, usize> = std::collections::HashMap::new(); - consume_m3u(cfg, input, lines.iter().cloned(), |item| { + consume_m3u(cfg, input, lines, |item| { + // keep the original sort order for groups and group the playlist items let key = Rc::clone(&item.header.borrow().group); match group_map.entry(key) { std::collections::hash_map::Entry::Vacant(v) => { - v.insert(idx); - idx += 1; + v.insert(sort_order_idx); sort_order.push(vec![item]); + sort_order_idx += 1; } std::collections::hash_map::Entry::Occupied(o) => { sort_order.get_mut(*o.get()).unwrap().push(item); } } }); - let mut grp_id = 0; let result: Vec = sort_order.drain(..).map(|channels| { + // create a group based on the first playlist item let channel = channels.first(); - let cluster = channel.map(|pli| pli.header.borrow().xtream_cluster).unwrap(); - let group_title = channel.map(|pli| Rc::clone(&pli.header.borrow().group)).unwrap(); + let (cluster, group_title) = channel.map(|pli| + (pli.header.borrow().xtream_cluster, Rc::clone(&pli.header.borrow().group))).unwrap(); grp_id += 1; PlaylistGroup { id: grp_id, xtream_cluster: cluster, title: Rc::clone(&group_title), channels } }).collect(); diff --git a/src/processing/playlist_processor.rs b/src/processing/playlist_processor.rs index 57cffdb51..c6d9b4d2b 100644 --- a/src/processing/playlist_processor.rs +++ b/src/processing/playlist_processor.rs @@ -25,8 +25,8 @@ use crate::processing::playlist_watch::process_group_watch; use crate::processing::xmltv_parser::flatten_tvguide; use crate::processing::xtream_processor::playlist_resolve_series; use crate::repository::playlist_repository::persist_playlist; -use crate::utils::download; use crate::utils::default_utils::default_as_default; +use crate::utils::download; fn is_valid(pli: &PlaylistItem, target: &ConfigTarget) -> bool { let provider = ValueProvider { pli: RefCell::new(pli) }; @@ -96,7 +96,7 @@ fn playlistitem_comparator(a: &PlaylistItem, b: &PlaylistItem, channel_sort: &Co } } } - }, + } None => { let ordering = value_a.partial_cmp(&value_b).unwrap(); match channel_sort.order { @@ -255,7 +255,7 @@ fn map_playlist_counter(target: &ConfigTarget, playlist: &[PlaylistGroup]) { for channel in &plg.channels { let provider = ValueProvider { pli: RefCell::new(channel) }; if counter.filter.filter(&provider, &mut mock_processor) { - let new_value = if counter.modifier == CounterModifier::Assign { + let new_value = if counter.modifier == CounterModifier::Assign { cntval.to_string() } else { let value = match channel.header.borrow_mut().get_field(&counter.field) { @@ -264,7 +264,7 @@ fn map_playlist_counter(target: &ConfigTarget, playlist: &[PlaylistGroup]) { }; if counter.modifier == CounterModifier::Suffix { format!("{}{}{}", value, counter.concat, cntval.to_string()) - } else { + } else { format!("{}{}{}", cntval.to_string(), counter.concat, value) } }; @@ -295,10 +295,11 @@ fn is_target_enabled(target: &ConfigTarget, user_targets: &ProcessTargets) -> bo async fn process_source(cfg: Arc, source_idx: usize, user_targets: Arc) -> (Vec, Vec) { let source = cfg.sources.get(source_idx).unwrap(); - let mut source_playlists = Vec::new(); - let enabled_inputs = source.inputs.iter().filter(|item| item.enabled).count(); let mut errors = vec![]; let mut stats = HashMap::::new(); + let mut source_playlists = Vec::new(); + let enabled_inputs = source.inputs.iter().filter(|item| item.enabled).count(); + // Downlod the sources for input in &source.inputs { let input_id = input.id; if is_input_enabled(enabled_inputs, input.enabled, input_id, &user_targets) { @@ -306,16 +307,16 @@ async fn process_source(cfg: Arc, source_idx: usize, user_targets: Arc

download::get_m3u_playlist(&cfg, input, &cfg.working_dir).await, InputType::Xtream => download::get_xtream_playlist(input, &cfg.working_dir).await, }; + // @TODO optmization dont hold tv_guide in memory, persist raw and later use sax parser to extract. let (tvguide, mut tvguide_errors) = if error_list.is_empty() { download::get_xmltv(&cfg, input, &cfg.working_dir).await } else { (None, vec![]) }; - error_list.drain(..).for_each(|err| errors.push(err)); - tvguide_errors.drain(..).for_each(|err| errors.push(err)); + errors.extend(error_list.drain(..).chain(tvguide_errors.drain(..))); let input_name = match &input.name { + Some(name_val) => name_val.as_str(), None => input.url.as_str(), - Some(name_val) => name_val.as_str() }; let group_count = playlistgroups.len(); let channel_count = playlistgroups.iter() @@ -334,19 +335,8 @@ async fn process_source(cfg: Arc, source_idx: usize, user_targets: Arc

, source_idx: usize, user_targets: Arc

{} - Err(mut err) => err.drain(..).for_each(|e| errors.push(e)) + Err(mut err) => errors.extend(err.drain(..)) } } } @@ -370,6 +360,22 @@ async fn process_source(cfg: Arc, source_idx: usize, user_targets: Arc

InputStats { + InputStats { + name: input_name.to_string(), + input_type: input_type, + error_count: error_count, + raw_stats: PlaylistStats { + group_count, + channel_count, + }, + processed_stats: PlaylistStats { + group_count: 0, + channel_count: 0, + }, + } +} + async fn process_sources(config: Arc, user_targets: Arc) -> (Vec, Vec) { let mut handle_list = vec![]; let thread_num = config.threads; @@ -390,10 +396,8 @@ async fn process_sources(config: Arc, user_targets: Arc) let (mut res_stats, mut res_errors) = System::new().block_on(async { process_source(cfg, index, usr_trgts).await }); - res_errors.drain(..) - .for_each(|err| shared_errors.lock().unwrap().push(err)); - res_stats.drain(..) - .for_each(|stat| shared_stats.lock().unwrap().push(stat)); + shared_errors.lock().unwrap().extend(res_errors.drain(..)); + shared_stats.lock().unwrap().extend(res_stats.drain(..)); }; handles.push(thread::spawn(process)); if handles.len() >= thread_num as usize { @@ -401,10 +405,8 @@ async fn process_sources(config: Arc, user_targets: Arc) } } else { let (mut res_stats, mut res_errors) = process_source(cfg, index, usr_trgts).await; - res_errors.drain(..) - .for_each(|err| shared_errors.lock().unwrap().push(err)); - res_stats.drain(..) - .for_each(|stat| shared_stats.lock().unwrap().push(stat)); + shared_errors.lock().unwrap().extend(res_errors.drain(..)); + shared_stats.lock().unwrap().extend(res_stats.drain(..)); } } for handle in handle_list { @@ -493,12 +495,11 @@ async fn process_playlist<'a>(playlists: &mut [FetchedPlaylist<'a>], let mut new_epg = vec![]; new_fetched_playlists.drain(..).for_each(|mut fp| { + // collect all epg_channel ids let epg_channel_ids: HashSet<_> = fp.playlistgroups.iter().flat_map(|g| &g.channels) .filter_map(|c| c.header.borrow().epg_channel_id.clone()).collect(); - fp.playlistgroups.drain(..).for_each(|group| { - new_playlist.push(group); - }); + new_playlist.extend(fp.playlistgroups.drain(..)); if !epg_channel_ids.is_empty() { if let Some(tv_guide) = fp.epg { debug!("found epg information for {}", &target.name); diff --git a/src/processing/xtream_parser.rs b/src/processing/xtream_parser.rs index 80dabba4a..96246bfbe 100644 --- a/src/processing/xtream_parser.rs +++ b/src/processing/xtream_parser.rs @@ -8,12 +8,11 @@ use crate::create_m3u_filter_error_result; use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; use crate::model::config::ConfigInput; use crate::model::playlist::{PlaylistGroup, PlaylistItem, PlaylistItemHeader, PlaylistItemType, XtreamCluster}; -use crate::model::xtream::{XtreamCategory, XtreamSeriesInfo, XtreamStream}; -use crate::utils::default_utils::default_as_empty_rc_str; +use crate::model::xtream::{XtreamCategory, XtreamSeriesInfo, XtreamSeriesInfoEpisode, XtreamStream}; -fn map_to_xtream_category(category: &Value) -> Result, M3uFilterError> { - match serde_json::from_value::>(category.to_owned()) { - Ok(category_list) => Ok(category_list), +fn map_to_xtream_category(categories: &Value) -> Result, M3uFilterError> { + match serde_json::from_value::>(categories.to_owned()) { + Ok(xtream_categories) => Ok(xtream_categories), Err(err) => { create_m3u_filter_error_result!(M3uFilterErrorKind::Notify, "Failed to process categories {}", &err) } @@ -29,6 +28,16 @@ fn map_to_xtream_streams(xtream_cluster: XtreamCluster, streams: &Value) -> Resu } } +fn create_xtream_series_info_url(url: &str, username: &str, password: &str, episode: &XtreamSeriesInfoEpisode) -> Rc { + if episode.direct_source.is_empty() { + let ext = episode.container_extension.clone(); + let stream_base_url = format!("{url}/series/{username}/{password}/{}.{ext}", episode.id); + Rc::new(stream_base_url) + } else { + Rc::new(episode.direct_source.clone()) + } +} + pub(crate) fn parse_xtream_series_info(info: &Value, group_title: &str, input: &ConfigInput) -> Result>, M3uFilterError> { let url = input.url.as_str(); let username = input.username.as_ref().map_or("", |v| v); @@ -39,33 +48,18 @@ pub(crate) fn parse_xtream_series_info(info: &Value, group_title: &str, input: & let result: Vec = series_info.episodes.values().flatten().map(|episode| PlaylistItem { header: RefCell::new(PlaylistItemHeader { - uuid: default_as_empty_rc_str(), - id: Rc::new(episode.id.clone()), - stream_id: Rc::new(episode.id.clone()), + id: Rc::new(episode.id.to_string()), name: Rc::new(episode.title.clone()), - chno: default_as_empty_rc_str(), logo: Rc::new(episode.info.movie_image.clone()), - logo_small: default_as_empty_rc_str(), group: Rc::new(group_title.to_string()), title: Rc::new(episode.title.clone()), - parent_code: default_as_empty_rc_str(), - audio_track: default_as_empty_rc_str(), - time_shift: default_as_empty_rc_str(), - rec: default_as_empty_rc_str(), - url: if episode.direct_source.is_empty() { - let ext = episode.container_extension.clone(); - let stream_base_url = format!("{}/series/{}/{}/{}.{}", url, username, password, episode.id.as_str(), ext); - Rc::new(stream_base_url) - } else { - Rc::new(episode.direct_source.clone()) - }, - epg_channel_id: None, + url: create_xtream_series_info_url(url, username, password, episode), item_type: PlaylistItemType::Series, xtream_cluster: XtreamCluster::Series, additional_properties: episode.get_additional_properties(&series_info), - series_fetched: false, category_id: 0, input_id: input.id, + ..Default::default() }) }).collect(); if result.is_empty() { Ok(None) } else { Ok(Some(result)) } @@ -76,69 +70,61 @@ pub(crate) fn parse_xtream_series_info(info: &Value, group_title: &str, input: & } } +fn create_xtream_url(xtream_cluster: XtreamCluster, url: &str, username: &str, password: &str, stream: &XtreamStream) -> Rc { + if stream.direct_source.is_empty() { + let stream_base_url = match xtream_cluster { + XtreamCluster::Live => format!("{}/live/{}/{}/{}.ts", url, username, password, &stream.get_stream_id()), + XtreamCluster::Video => { + let ext = stream.container_extension.as_ref().map_or("mp4", |e| e.as_str()); + format!("{}/movie/{}/{}/{}.{}", url, username, password, &stream.get_stream_id(), ext) + } + XtreamCluster::Series => + format!("{}/player_api.php?username={}&password={}&action=get_series_info&series_id={}", + url, username, password, &stream.get_stream_id()) + }; + Rc::new(stream_base_url) + } else { + Rc::clone(&stream.direct_source) + } +} pub(crate) fn parse_xtream(input: &ConfigInput, xtream_cluster: XtreamCluster, - category: &Value, + categories: &Value, streams: &Value) -> Result>, M3uFilterError> { - match map_to_xtream_category(category) { - Ok(mut categories) => { + match map_to_xtream_category(categories) { + Ok(mut xtream_categories) => { let input_id = input.id; let url = input.url.as_str(); let username = input.username.as_ref().map_or("", |v| v); let password = input.password.as_ref().map_or("", |v| v); return match map_to_xtream_streams(xtream_cluster, streams) { - Ok(streams) => { + Ok(xtream_streams) => { let group_map: HashMap::, RefCell> = - categories.drain(..).map(|category| + xtream_categories.drain(..).map(|category| (Rc::clone(&category.category_id), RefCell::new(category)) ).collect(); - for stream in streams { + for stream in xtream_streams { if let Some(group) = group_map.get(&stream.category_id) { let mut grp = group.borrow_mut(); let category_name = &grp.category_name; let item = PlaylistItem { header: RefCell::new(PlaylistItemHeader { - uuid: default_as_empty_rc_str(), - id: Rc::new(stream.get_stream_id()), - stream_id: Rc::new(stream.get_stream_id()), + id: Rc::new(stream.get_stream_id().to_string()), name: Rc::clone(&stream.name), - chno: default_as_empty_rc_str(), logo: Rc::clone(&stream.stream_icon), - logo_small: default_as_empty_rc_str(), group: Rc::clone(category_name), title: Rc::clone(&stream.name), - parent_code: default_as_empty_rc_str(), - audio_track: default_as_empty_rc_str(), - time_shift: default_as_empty_rc_str(), - rec: default_as_empty_rc_str(), - url: if stream.direct_source.is_empty() { - let stream_base_url = match xtream_cluster { - XtreamCluster::Live => format!("{}/live/{}/{}/{}.ts", url, username, password, &stream.get_stream_id()), - XtreamCluster::Video => { - let ext = stream.container_extension.as_ref().map_or("mp4", |e| e.as_str()); - format!("{}/movie/{}/{}/{}.{}", url, username, password, &stream.get_stream_id(), ext) - } - XtreamCluster::Series => - format!("{}/player_api.php?username={}&password={}&action=get_series_info&series_id={}", - url, username, password, &stream.get_stream_id()) - }; - Rc::new(stream_base_url) - } else { - Rc::clone(&stream.direct_source) - }, + url: create_xtream_url(xtream_cluster, url, username, password, &stream), epg_channel_id: stream.epg_channel_id.clone(), - item_type: match xtream_cluster { - XtreamCluster::Live => PlaylistItemType::Live, - XtreamCluster::Video => PlaylistItemType::Movie, - XtreamCluster::Series => PlaylistItemType::SeriesInfo, - }, + item_type: PlaylistItemType::from(xtream_cluster), xtream_cluster, additional_properties: stream.get_additional_properties(), series_fetched: false, category_id: 0, input_id, + ..Default::default() }), }; grp.add(item); @@ -160,4 +146,5 @@ pub(crate) fn parse_xtream(input: &ConfigInput, } Err(err) => Err(err) } -} \ No newline at end of file +} + diff --git a/src/repository/bplustree.rs b/src/repository/bplustree.rs deleted file mode 100644 index 03127ae81..000000000 --- a/src/repository/bplustree.rs +++ /dev/null @@ -1,181 +0,0 @@ -use std::fs::File; -use std::io::{self}; -use flate2::Compression; -use flate2::write::{GzEncoder}; -use flate2::read::{GzDecoder}; - -use serde::{Deserialize, Serialize}; - -const T: usize = 3; // Minimum degree (T), meaning each node can contain at most 2*T - 1 keys - -#[derive(Serialize, Deserialize, Debug, Clone)] -struct KeyValue { - key: K, - value: V, -} - -#[derive(Serialize, Deserialize, Debug, Clone)] -struct BPlusTreeNode { - keys: Vec, - children: Vec>>, - is_leaf: bool, - values: Vec>, // only used in leaf nodes -} - -impl BPlusTreeNode -where - K: Ord + Serialize + for<'de> Deserialize<'de> + Clone, - V: Serialize + for<'de> Deserialize<'de> + Clone, -{ - fn new(is_leaf: bool) -> Self { - BPlusTreeNode { - keys: vec![], - children: vec![], - is_leaf, - values: vec![], - } - } - - fn insert_non_full(&mut self, key: K, value: V) { - let pos = self.keys.binary_search(&key).unwrap_or_else(|pos| pos); - - if self.is_leaf { - self.keys.insert(pos, key); - self.values.insert(pos, Some(value)); - } else { - let mut idx = pos; - if self.children[pos].as_mut().unwrap().keys.len() == 2 * T - 1 { - self.split_child(pos); - if key > self.keys[pos] { - idx = pos + 1; - } - } - self.children[idx].as_mut().unwrap().insert_non_full(key, value); - } - } - - fn split_child(&mut self, pos: usize) { - let t = T - 1; - let mut new_node = BPlusTreeNode::new(self.children[pos].as_ref().unwrap().is_leaf); - let mut old_node = self.children[pos].take().unwrap(); - - self.keys.insert(pos, old_node.keys.remove(t)); - self.children.insert(pos + 1, Some(new_node.clone())); - - if old_node.is_leaf { - new_node.keys = old_node.keys.split_off(t); - new_node.values = old_node.values.split_off(t); - } else { - new_node.keys = old_node.keys.split_off(t + 1); - new_node.children = old_node.children.split_off(t + 1); - } - - self.children[pos] = Some(old_node); - self.children[pos + 1] = Some(new_node); - } -} - -#[derive(Serialize, Deserialize, Debug, Clone)] -struct BPlusTree { - root: BPlusTreeNode, -} - -impl BPlusTree -where - K: Ord + Serialize + for<'de> Deserialize<'de> + Clone, - V: Serialize + for<'de> Deserialize<'de> + Clone, -{ - pub(crate) fn new() -> Self { - BPlusTree { - root: BPlusTreeNode::new(true), - } - } - - pub(crate) fn insert(&mut self, key: K, value: V) { - if self.root.keys.len() == 2 * T - 1 { - let mut new_root = BPlusTreeNode::new(false); - new_root.children.push(Some(self.root.clone())); - new_root.split_child(0); - self.root = new_root; - } - self.root.insert_non_full(key, value); - } - - pub(crate) fn query(&self, key: &K) -> Option { - let mut node = &self.root; - while !node.is_leaf { - let pos = match node.keys.binary_search(key) { - Ok(pos) => return node.values[pos].clone(), - Err(pos) => pos, - }; - node = node.children[pos].as_ref().unwrap(); - } - - match node.keys.binary_search(key) { - Ok(pos) => node.values[pos].clone(), - Err(_) => None, - } - } - - pub(crate) fn serialize_to_file(&self, filename: &str) -> io::Result<()> { - let file = File::create(filename)?; - let encoder = GzEncoder::new(file, Compression::default()); - match bincode::serialize_into(encoder, &self) { - Ok(()) => Ok(()), - Err(e) => { - println!("Failed to write bplustree to disk {e}"); - Ok(()) - } - } - } - - // If file exists the file is deserialized, otherwise an empty tree is returned - pub(crate) fn deserialize_from_file(filename: &str) -> Self { - match File::open(filename) { - Ok(file) => { - let decoder = GzDecoder::new(file); - let tree: BPlusTree = bincode::deserialize_from(decoder).unwrap(); - tree - } - Err(_) => BPlusTree::new() - } - } -} - -#[cfg(test)] -mod tests { - use std::io; - use serde::{Deserialize, Serialize}; - use crate::repository::bplustree::{BPlusTree}; - - - // Example usage with a simple struct - #[derive(Serialize, Deserialize, Debug, Clone)] - struct Value { - id: u32, - data: String, - } - - #[test] - fn insert_test() -> io::Result<()> { - let mut tree = BPlusTree::new(); - tree.insert("abc".to_string(), Value { id: 16, data: "one".to_string() }); - tree.insert("def".to_string(), Value { id: 32, data: "two".to_string() }); - tree.insert("ghi".to_string(), Value { id: 64, data: "three".to_string() }); - - // Serialize the tree to a file - tree.serialize_to_file("/tmp/tree.bin")?; - - // Deserialize the tree from the file - let tree = BPlusTree::::deserialize_from_file("/tmp/tree.bin"); - - // Query the tree - if let Some(value) = tree.query(&("ghi".to_string())) { - println!("Found: {:?}", value); - } else { - println!("Not found"); - } - - Ok(()) - } -} \ No newline at end of file diff --git a/src/repository/epg_repository.rs b/src/repository/epg_repository.rs index 211930b59..528497d71 100644 --- a/src/repository/epg_repository.rs +++ b/src/repository/epg_repository.rs @@ -42,7 +42,7 @@ fn epg_write_file(target: &ConfigTarget, epg: &Epg, path: &Path) -> Result<(), M Ok(()) } -pub(crate) fn epg_write(target: &ConfigTarget, cfg: &Config, epg: Option<&Epg>, output: &TargetOutput) -> Result<(), M3uFilterError> { +pub(crate) fn epg_write(target: &ConfigTarget, cfg: &Config, target_path: &Path, epg: Option<&Epg>, output: &TargetOutput) -> Result<(), M3uFilterError> { if let Some(epg_data) = epg { match &output.target { TargetType::M3u => { @@ -51,7 +51,7 @@ pub(crate) fn epg_write(target: &ConfigTarget, cfg: &Config, epg: Option<&Epg>, M3uFilterErrorKind::Notify, format!("write epg for target {} failed: No filename set", target.name))); } - if let Some(path) = m3u_get_epg_file_path(cfg, target) { + if let Some(path) = m3u_get_epg_file_path(target_path) { if log_enabled!(Level::Debug) { debug!("writing m3u epg to {}", path.to_str().unwrap_or("?")); } diff --git a/src/repository/index_record.rs b/src/repository/index_record.rs index 1bb0e7808..c4cf37ac1 100644 --- a/src/repository/index_record.rs +++ b/src/repository/index_record.rs @@ -39,23 +39,6 @@ impl IndexRecord { } } - // pub fn from_bytes(bytes: &[u8], cursor: &mut usize) -> Result { - // if let Ok(index_bytes) = bytes[*cursor..*cursor + 4].try_into() { - // *cursor += 4; - // if let Ok(size_bytes) = bytes[*cursor..*cursor + 2].try_into() { - // *cursor += 2; - // let index = u32::from_le_bytes(index_bytes); - // let size = u16::from_le_bytes(size_bytes); - // return Ok(IndexRecord { index, size }); - // } - // } - // Err(Error::new(ErrorKind::Other, "Failed to read index")) - // } - - // pub fn as_bytes(&self) -> [u8; 6] { - // IndexRecord::to_bytes(self.index, self.size) - // } - pub fn to_bytes(left: u32, right: u32) -> [u8; 8] { let left_bytes: [u8; 4] = left.to_le_bytes(); let right_bytes: [u8; 4] = right.to_le_bytes(); diff --git a/src/repository/indexed_document_writer.rs b/src/repository/indexed_document_writer.rs index 289bdca7f..b05d7439f 100644 --- a/src/repository/indexed_document_writer.rs +++ b/src/repository/indexed_document_writer.rs @@ -49,7 +49,7 @@ impl IndexedDocumentWriter { Self::new_with_mode(main_path, index_path, true) } - pub fn write_doc(&mut self, doc: &T) -> Result<(u32, u32), Error> + pub fn write_doc(&mut self, doc_id: u32, doc: &T) -> Result<(u32, u32), Error> where T: ?Sized + serde::Serialize { let current_main_index = self.main_offset; @@ -58,7 +58,7 @@ impl IndexedDocumentWriter { match file_utils::check_write(&self.main_file.write_all(&encoded)) { Ok(()) => { let bytes_written = u32::try_from(encoded.len()).map_err(|err| Error::new(ErrorKind::Other, err))?; - let combined_bytes = IndexRecord::to_bytes(self.main_offset, bytes_written); + let combined_bytes = IndexRecord::to_bytes(doc_id, bytes_written); if let Err(err) = file_utils::check_write(&self.index_file.write_all(&combined_bytes)) { return Err(Error::new(ErrorKind::Other, format!("failed to write document: {} - {}", self.index_path.to_str().unwrap(), err))); } diff --git a/src/repository/m3u_repository.rs b/src/repository/m3u_repository.rs index 4f96daa97..214c73b4d 100644 --- a/src/repository/m3u_repository.rs +++ b/src/repository/m3u_repository.rs @@ -1,7 +1,6 @@ use std::fs::File; use std::io::{BufWriter, Error, ErrorKind, Write}; use std::path::{Path, PathBuf}; -use std::rc::Rc; use log::error; @@ -14,6 +13,7 @@ use crate::model::playlist::{M3uPlaylistItem, PlaylistGroup, PlaylistItem, Playl use crate::repository::index_record::IndexRecord; use crate::repository::indexed_document_reader::{IndexedDocumentReader, read_indexed_item}; use crate::repository::indexed_document_writer::IndexedDocumentWriter; +use crate::repository::storage::{ensure_target_storage_path}; use crate::utils::file_utils; macro_rules! cant_write_result { @@ -22,12 +22,12 @@ macro_rules! cant_write_result { } } -fn m3u_get_base_file_path(cfg: &Config, target: &ConfigTarget) -> Option { - file_utils::get_file_path(&cfg.working_dir, Some(PathBuf::from(format!("m3u_{}.db", target.name.replace(' ', "_").as_str())))) +fn m3u_get_base_file_path(target_path: &Path) -> Option { + Some(target_path.join(PathBuf::from("m3u.db"))) } -pub(crate) fn m3u_get_file_paths(cfg: &Config, target: &ConfigTarget) -> Option<(PathBuf, PathBuf)> { - match m3u_get_base_file_path(cfg, target) { +pub(crate) fn m3u_get_file_paths(target_path: &Path) -> Option<(PathBuf, PathBuf)> { + match m3u_get_base_file_path(target_path) { Some(m3u_path) => { let extension = m3u_path.extension().map(|ext| format!("{}_", ext.to_str().unwrap_or(""))); let index_path = m3u_path.with_extension(format!("{}idx", &extension.unwrap_or(String::new()))); @@ -37,44 +37,47 @@ pub(crate) fn m3u_get_file_paths(cfg: &Config, target: &ConfigTarget) -> Option< } } -pub(crate) fn m3u_get_epg_file_path(cfg: &Config, target: &ConfigTarget) -> Option { - m3u_get_base_file_path(cfg, target) +pub(crate) fn m3u_get_epg_file_path(target_path: &Path) -> Option { + m3u_get_base_file_path(target_path) .map(|path| file_utils::add_prefix_to_filename(&path, "epg_", Some("xml"))) } -pub(crate) fn m3u_write_playlist(target: &ConfigTarget, cfg: &Config, new_playlist: &[PlaylistGroup]) -> Result<(), M3uFilterError> { +fn persist_m3u_playlist_as_text(target: &ConfigTarget, cfg: &Config, m3u_playlist: &Vec) { + if let Some(filename) = target.get_m3u_filename() { + if let Some(m3u_filename) = file_utils::get_file_path(&cfg.working_dir, Some(PathBuf::from(filename))) { + match File::create(&m3u_filename) { + Ok(file) => { + let mut buf_writer = BufWriter::new(file); + let _ = buf_writer.write("#EXTM3U\n".as_bytes()); + for m3u in m3u_playlist { + let _ = buf_writer.write(m3u.to_m3u(target, None).as_bytes()); + let _ = buf_writer.write("\n".as_bytes()); + } + } + Err(_) => { + error!("Can't write m3u plain playlist {}", &m3u_filename.to_str().unwrap()); + } + } + } + } +} + +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(cfg, target) { + 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::>(); - if let Some(filename) = target.get_m3u_filename() { - if let Some(m3u_filename) = file_utils::get_file_path(&cfg.working_dir, Some(PathBuf::from(filename))) { - match File::create(&m3u_filename) { - Ok(file) => { - let mut buf_writer = BufWriter::new(file); - let _ = buf_writer.write("#EXTM3U\n".as_bytes()); - for m3u in &m3u_playlist { - let _ = buf_writer.write(m3u.to_m3u(target, None).as_bytes()); - let _ = buf_writer.write("\n".as_bytes()); - } - } - Err(_) => { - error!("Can't write m3u plain playlist {}", &m3u_filename.to_str().unwrap()); - } - } - } - } + persist_m3u_playlist_as_text(target, cfg, &m3u_playlist); + match IndexedDocumentWriter::new(m3u_path.clone(), idx_path) { Ok(mut writer) => { - let mut stream_id: u32 = 1; - for mut m3u in m3u_playlist { - m3u.stream_id = Rc::new(stream_id.to_string()); - match writer.write_doc(&m3u) { - Ok(_) => stream_id += 1, + for m3u in m3u_playlist { + match writer.write_doc(m3u.virtual_id, &m3u) { + Ok(_) => {}, Err(err) => return cant_write_result!(&m3u_path, err) } } @@ -87,36 +90,43 @@ pub(crate) fn m3u_write_playlist(target: &ConfigTarget, cfg: &Config, new_playli } pub(crate) fn m3u_load_rewrite_playlist(cfg: &Config, target: &ConfigTarget, user: &ProxyUserCredentials) -> Option { - if let Some((m3u_path, idx_path)) = m3u_get_file_paths(cfg, target) { - match IndexedDocumentReader::::new(&m3u_path, &idx_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 => { - let stream_id = Rc::clone(&m3u_pli.stream_id); - result.push(m3u_pli.to_m3u(target, Some(format!("{url}/{stream_id}").as_str()))); - } - ProxyType::Redirect => { - result.push(m3u_pli.to_m3u(target, None)); + match ensure_target_storage_path(cfg, target.name.as_str()) { + Ok(target_path) => { + if let Some((m3u_path, idx_path)) = m3u_get_file_paths(&target_path) { + match IndexedDocumentReader::::new(&m3u_path, &idx_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 => { + let stream_id = m3u_pli.virtual_id; + result.push(m3u_pli.to_m3u(target, Some(format!("{url}/{stream_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); + } } + } else { + error!("Could not open files for target {}", &target.name); } - Err(err) => { - error!("Could not deserialize file {} - {}", &m3u_path.to_str().unwrap(), err); - } + }, + Err(err) => { + error!("Could not find storage path for target {} - {}", target.name.as_str(), err); } - } else { - error!("Could not open files for target {}", &target.name); } None } diff --git a/src/repository/mod.rs b/src/repository/mod.rs index a1f08200b..8ca12a2fa 100644 --- a/src/repository/mod.rs +++ b/src/repository/mod.rs @@ -2,10 +2,10 @@ pub(crate) mod playlist_repository; pub(crate) mod m3u_repository; pub(crate) mod xtream_repository; pub(crate) mod epg_repository; -pub(crate)mod kodi_repository; - -mod bplustree; +pub(crate) mod kodi_repository; +pub(crate) mod storage; mod index_record; mod indexed_document_writer; -mod indexed_document_reader; \ No newline at end of file +mod indexed_document_reader; +mod target_id_mapping_record; \ No newline at end of file diff --git a/src/repository/playlist_repository.rs b/src/repository/playlist_repository.rs index 3a820f793..51662fc53 100644 --- a/src/repository/playlist_repository.rs +++ b/src/repository/playlist_repository.rs @@ -1,31 +1,59 @@ -use crate::m3u_filter_error::M3uFilterError; +use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; use crate::model::config::{Config, ConfigTarget, TargetType}; use crate::model::playlist::PlaylistGroup; use crate::model::xmltv::Epg; use crate::repository::epg_repository::epg_write; use crate::repository::kodi_repository::kodi_write_strm_playlist; use crate::repository::m3u_repository::m3u_write_playlist; +use crate::repository::storage::{ensure_target_storage_path, get_target_id_mapping_file}; +use crate::repository::target_id_mapping_record::TargetIdMapping; use crate::repository::xtream_repository::xtream_write_playlist; pub(crate) fn persist_playlist(playlist: &mut [PlaylistGroup], epg: Option<&Epg>, target: &ConfigTarget, cfg: &Config) -> Result<(), Vec> { let mut errors = vec![]; - for output in &target.output { - match match output.target { - TargetType::M3u => m3u_write_playlist(target, cfg, playlist), - TargetType::Xtream => xtream_write_playlist(target, cfg, playlist), - TargetType::Strm => kodi_write_strm_playlist(target, cfg, playlist, &output.filename), - } { - Ok(()) => { - if !playlist.is_empty() { - match epg_write(target, cfg, epg, output) { - Ok(()) => {} - Err(err) => errors.push(err) + + // TODO get previous virtual-ids and match them to the playlist items + match ensure_target_storage_path(cfg, target.name.as_str()) { + Ok(target_path) => { + let mut target_id_mapping = TargetIdMapping::from_path(&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.id.parse::() { + Ok(provider_id) => { + let uuid = header.get_uuid(); + header.virtual_id = target_id_mapping.insert_entry(provider_id, **uuid); + } + Err(err) => { + errors.push(M3uFilterError::new(M3uFilterErrorKind::Info, err.to_string())); + } } } } - Err(err) => errors.push(err) + + 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, &target_path, playlist), + TargetType::Strm => kodi_write_strm_playlist(target, cfg, playlist, &output.filename), + } { + Ok(()) => { + if let Err(err) = target_id_mapping.to_path(&target_path) { + 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) + } + } } + 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 new file mode 100644 index 000000000..62e899703 --- /dev/null +++ b/src/repository/storage.rs @@ -0,0 +1,30 @@ +use std::path::{Path, PathBuf}; +use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; +use crate::model::config::Config; +use crate::utils::file_utils; + +pub(crate) fn hash_string(url: &str) -> [u8; 32] { + let hash = blake3::hash(url.as_bytes()); + hash.into() // Konvertiere den Hash in ein Array mit fester Größe +} + +pub(crate) fn get_target_id_mapping_file(target_path: &Path) -> PathBuf { + target_path.join(PathBuf::from("id_mapping.db")) +} + +pub(crate) fn ensure_target_storage_path(cfg: &Config, target_name: &str) -> Result { + if let Some(path) = get_target_storage_path(cfg, target_name) { + if std::fs::create_dir_all(&path).is_err() { + let msg = format!("Failed to save target data, can't create directory {}", &path.to_str().unwrap()); + return Err(M3uFilterError::new(M3uFilterErrorKind::Notify, msg)); + } + Ok(path) + } else { + let msg = format!("Failed to save target data, can't create directory for target {target_name}"); + Err(M3uFilterError::new(M3uFilterErrorKind::Notify, msg)) + } +} + +pub(crate) fn get_target_storage_path(cfg: &Config, target_name: &str) -> Option { + file_utils::get_file_path(&cfg.working_dir, Some(std::path::PathBuf::from(target_name.replace(' ', "_")))) +} diff --git a/src/repository/target_id_mapping_record.rs b/src/repository/target_id_mapping_record.rs new file mode 100644 index 000000000..7e638ddae --- /dev/null +++ b/src/repository/target_id_mapping_record.rs @@ -0,0 +1,141 @@ +use std::cmp::max; +use std::collections::BTreeMap; +use std::fs::File; +use std::io::{Error, Read, Seek, SeekFrom, Write}; +use std::path::Path; +use std::rc::Rc; + +use crate::utils::file_utils; + +/** +This file contains the provider id, the virtual id, and an uuid + + */ +pub(in crate::repository) struct TargetIdMappingRecord { + pub provider_id: u32, + pub virtual_id: u32, + pub uuid: [u8; 32], +} + +impl TargetIdMappingRecord { + fn to_bytes(&self) -> [u8; 40] { + let provider_id_bytes: [u8; 4] = self.provider_id.to_le_bytes(); + let virtual_id_bytes: [u8; 4] = self.virtual_id.to_le_bytes(); + let mut combined_bytes: [u8; 40] = [0; 40]; + combined_bytes[..4].copy_from_slice(&provider_id_bytes); + combined_bytes[4..].copy_from_slice(&virtual_id_bytes); + combined_bytes[8..].copy_from_slice(&self.uuid); + combined_bytes + } +} + +pub(in crate::repository) struct TargetIdMapping { + dirty: bool, + virtual_id_counter: u32, + by_provider_id: BTreeMap>, + by_virtual_id: BTreeMap>, + by_uuid: BTreeMap<[u8; 32], Rc>, + records: Vec>, +} + +impl TargetIdMapping { + pub(crate) fn insert_entry(&mut self, provider_id: u32, uuid: [u8; 32]) -> u32 { + self.dirty = true; + self.virtual_id_counter += 1; + self.insert(TargetIdMappingRecord { provider_id, virtual_id: self.virtual_id_counter, uuid }); + self.virtual_id_counter + } + + fn new(records: Vec>) -> Self { + let mut by_provider_id: BTreeMap> = BTreeMap::new(); + let mut by_virtual_id: BTreeMap> = BTreeMap::new(); + let mut by_uuid: BTreeMap<[u8; 32], Rc> = BTreeMap::new(); + let mut virtual_id_counter: u32 = 0; + for record in &records { + by_provider_id.insert(record.provider_id, Rc::clone(record)); + by_virtual_id.insert(record.virtual_id, Rc::clone(record)); + by_uuid.insert(record.uuid, Rc::clone(record)); + virtual_id_counter = max(record.virtual_id, virtual_id_counter); + } + TargetIdMapping { + dirty: false, + virtual_id_counter, + by_provider_id, + by_virtual_id, + by_uuid, + records, + } + } + + fn insert(&mut self, record: TargetIdMappingRecord) { + let provider_id = record.provider_id; + let virtual_id = record.virtual_id; + let uuid = record.uuid; + let shared_record = Rc::new(record); + self.by_provider_id.insert(provider_id, Rc::clone(&shared_record)); + self.by_virtual_id.insert(virtual_id, Rc::clone(&shared_record)); + self.by_uuid.insert(uuid, Rc::clone(&shared_record)); + self.records.push(shared_record); + } + + fn get_by_provider_id(&self, provider_id: u32) -> Option<&Rc> { + self.by_provider_id.get(&provider_id) + } + + fn get_by_virtual_id(&self, virtual_id: u32) -> Option<&Rc> { + self.by_virtual_id.get(&virtual_id) + } + + fn get_by_uuid(&self, uuid: &[u8; 32]) -> Option<&Rc> { + self.by_uuid.get(uuid) + } + + pub fn to_file(&mut self, file: &mut File) -> Result<(), Error> { + for record in &self.records { + let bytes = record.to_bytes(); + if let Err(err) = file.write_all(&bytes) { + return Err(err); + } + } + self.dirty = false; + Ok(()) + } + + pub fn from_path(path: &Path) -> Self { + match file_utils::open_file_append(path, false) { + Ok(mut file) => TargetIdMapping::from_file(&mut file), + _ => TargetIdMapping::new(vec![]) + } + } + + pub fn to_path(&mut self, path: &Path) -> Result<(), Error> { + match file_utils::open_file_append(path, false) { + Ok(mut file) => self.to_file(&mut file), + Err(err) => Err(err) + } + } + + pub fn from_file(file: &mut File) -> Self { + let mut records = vec![]; + if let Ok(_) = file.seek(SeekFrom::Start(0)) { + loop { + let mut provider_id_bytes = [0u8; 4]; + let mut virtual_id_bytes = [0u8; 4]; + let mut uuid = [0u8; 32]; + if let Err(_) = file.read_exact(&mut provider_id_bytes) { + break; + } + if let Err(_) = file.read_exact(&mut virtual_id_bytes) { + break; + } + if let Err(_) = file.read_exact(&mut uuid) { + break; + } + let provider_id = u32::from_le_bytes(provider_id_bytes); + let virtual_id = u32::from_le_bytes(virtual_id_bytes); + records.push(Rc::new(TargetIdMappingRecord { provider_id, virtual_id, uuid })); + } + } + TargetIdMapping::new(records) + } +} \ No newline at end of file diff --git a/src/repository/xtream_repository.rs b/src/repository/xtream_repository.rs index a3c362712..7fbed66c8 100644 --- a/src/repository/xtream_repository.rs +++ b/src/repository/xtream_repository.rs @@ -17,6 +17,7 @@ use crate::processing::m3u_parser::extract_id_from_url; use crate::repository::index_record::IndexRecord; use crate::repository::indexed_document_reader::{IndexedDocumentReader, read_indexed_item}; use crate::repository::indexed_document_writer::IndexedDocumentWriter; +use crate::repository::storage::get_target_storage_path; use crate::utils::file_utils; use crate::utils::json_utils::{json_iter_array, json_write_documents_to_file}; @@ -100,8 +101,8 @@ fn write_playlist_to_file(storage_path: &Path, stream_id: &mut u32, cluster: Xtr Ok(mut writer) => { for pli in playlist.iter_mut() { if let Ok(mut xtream) = pli.to_xtream() { - xtream.stream_id = *stream_id; - match writer.write_doc(&xtream) { + xtream.virtual_id = *stream_id; + match writer.write_doc(xtream.virtual_id, &xtream) { Ok(_) => *stream_id += 1, Err(err) => return cant_write_result!(&xtream_path, err) } @@ -193,7 +194,10 @@ fn load_old_category_ids(path: &Path) -> (u32, HashMap) { } pub(crate) fn xtream_get_storage_path(cfg: &Config, target_name: &str) -> Option { - file_utils::get_file_path(&cfg.working_dir, Some(std::path::PathBuf::from(target_name.replace(' ', "_")))) + match get_target_storage_path(cfg, target_name) { + Some(target_path) => Some(target_path.join(std::path::PathBuf::from("xtream"))), + None => None, + } } pub(crate) fn xtream_get_epg_file_path(path: &Path) -> PathBuf { @@ -211,7 +215,7 @@ pub(crate) fn xtream_get_file_paths(storage_path: &Path, cluster: XtreamCluster) (xtream_path, index_path) } -pub(crate) fn xtream_write_playlist(target: &ConfigTarget, cfg: &Config, playlist: &mut [PlaylistGroup]) -> Result<(), M3uFilterError> { +pub(crate) fn xtream_write_playlist(target: &ConfigTarget, cfg: &Config, target_path: &Path, playlist: &mut [PlaylistGroup]) -> Result<(), M3uFilterError> { match ensure_xtream_storage_path(cfg, target.name.replace(' ', "_").as_str()) { Ok(path) => { let mut cat_live_col = vec![]; @@ -221,7 +225,7 @@ pub(crate) fn xtream_write_playlist(target: &ConfigTarget, cfg: &Config, playlis let mut series_col = vec![]; let mut vod_col = vec![]; let mut errors = Vec::new(); - +// TODO.EUZU old catgeory ids // preserve category_ids let (max_cat_id, existing_cat_ids) = load_old_category_ids(&path); let mut cat_id_counter = max_cat_id; @@ -473,7 +477,7 @@ pub(crate) fn xtream_write_series_info(config: &Config, target_name: &str, 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(content) { + match writer.write_doc(series_id, content) { Ok((_, index_offset)) => { let series_id_index_mapping_path = xtream_get_series_id_series_info_mapping_file_path(&storage_path); IndexRecord::to_file(&series_id_index_mapping_path, series_id, index_offset, true)?; diff --git a/src/test.rs b/src/test.rs index 4cd838d24..cd7de9006 100644 --- a/src/test.rs +++ b/src/test.rs @@ -1,35 +1,14 @@ #[cfg(test)] mod tests { use std::cell::RefCell; - use std::rc::Rc; use regex::Regex; use crate::filter::{get_filter, MockValueProcessor, ValueProvider}; - use crate::model::playlist::{PlaylistItem, PlaylistItemHeader, PlaylistItemType, XtreamCluster}; + use crate::model::playlist::{PlaylistItem, PlaylistItemHeader}; fn create_mock_pli(name: &str, group: &str) -> PlaylistItem { PlaylistItem { header: RefCell::new(PlaylistItemHeader { - uuid: Rc::new("".to_string()), - stream_id: Rc::new("".to_string()), - id: Rc::new("".to_string()), - name: Rc::new(name.to_string()), - chno: Rc::new("".to_string()), - logo: Rc::new("".to_string()), - logo_small: Rc::new("".to_string()), - group: Rc::new(group.to_string()), - title: Rc::new("".to_string()), - parent_code: Rc::new("".to_string()), - audio_track: Rc::new("".to_string()), - time_shift: Rc::new("".to_string()), - rec: Rc::new("".to_string()), - url: Rc::new("".to_string()), - epg_channel_id: None, - xtream_cluster: XtreamCluster::Live, - additional_properties: None, - item_type: PlaylistItemType::Live, - series_fetched: false, - category_id: 0, - input_id: 0, + ..Default::default() }) } } diff --git a/src/utils/default_utils.rs b/src/utils/default_utils.rs index 758436c94..e51dbec7f 100644 --- a/src/utils/default_utils.rs +++ b/src/utils/default_utils.rs @@ -1,7 +1,6 @@ use std::collections::HashMap; use std::rc::Rc; use crate::model::config::ProcessingOrder; -use crate::model::playlist::{PlaylistItemType, XtreamCluster}; pub(crate) fn default_as_true() -> bool { true } @@ -23,8 +22,6 @@ pub(crate) fn default_as_empty_list() -> Vec { vec![] } pub(crate) fn default_as_two_u16() -> u16 { 2 } -pub(crate) fn default_playlist_item_type() -> PlaylistItemType { PlaylistItemType::Live } pub(crate) fn default_as_zero_u32() -> u32 { 0 } pub(crate) fn default_as_zero_u16() -> u16 { 0 } -pub(crate) fn default_stream_cluster() -> XtreamCluster { XtreamCluster::Live } diff --git a/src/utils/download.rs b/src/utils/download.rs index 9a0bf2358..a5ddff2de 100644 --- a/src/utils/download.rs +++ b/src/utils/download.rs @@ -30,8 +30,7 @@ pub(crate) async fn get_m3u_playlist(cfg: &Config, input: &ConfigInput, working_ let persist_file_path = prepare_file_path(input.persist.as_ref(), working_dir, ""); match request_utils::get_input_text_content(input, working_dir, &url, persist_file_path).await { Ok(text) => { - let lines = text.lines().map(String::from).collect::>(); - (m3u_parser::parse_m3u(cfg, input, &lines), vec![]) + (m3u_parser::parse_m3u(cfg, input, text.lines()), vec![]) } Err(err) => (vec![], vec![err]) } @@ -57,7 +56,7 @@ pub(crate) async fn get_xtream_playlist_series<'a>(fpl: &mut FetchedPlaylist<'a> match parse_xtream_series_info(&series_content, pli.header.borrow().group.as_str(), input) { Ok(series_info) => { if let Some(mut series) = series_info { - series.drain(..).for_each(|item| group_series.push(item)); + group_series.extend(series.drain(..)); } } Err(err) => errors.push(err), @@ -108,7 +107,7 @@ const ACTIONS: [(XtreamCluster, &str, &str); 3] = [ (XtreamCluster::Series, "get_series_categories", "get_series")]; pub(crate) async fn get_xtream_playlist(input: &ConfigInput, working_dir: &String) -> (Vec, Vec) { - let mut playlist: Vec = Vec::new(); + let mut playlist_groups: Vec = Vec::new(); let username = input.username.as_ref().map_or("", |v| v); let password = input.password.as_ref().map_or("", |v| v); let base_url = format!("{}/player_api.php?username={}&password={}", input.url, username, password); @@ -123,35 +122,37 @@ pub(crate) async fn get_xtream_playlist(input: &ConfigInput, working_dir: &Strin let category_file_path = prepare_file_path(input.persist.as_ref(), working_dir, format!("{category}_").as_str()); let stream_file_path = prepare_file_path(input.persist.as_ref(), working_dir, format!("{stream}_").as_str()); - match request_utils::get_input_json_content(input, category_url.as_str(), category_file_path).await { - Ok(category_content) => { - match request_utils::get_input_json_content(input, stream_url.as_str(), stream_file_path).await { - Ok(stream_content) => { - match xtream_parser::parse_xtream(input, - *xtream_cluster, - &category_content, - &stream_content) { - Ok(sub_playlist_opt) => { - if let Some(mut sub_playlist) = sub_playlist_opt { - sub_playlist.drain(..).for_each(|group| playlist.push(group)); - } - } - Err(err) => errors.push(err) + match futures::join!( + request_utils::get_input_json_content(input, category_url.as_str(), category_file_path), + request_utils::get_input_json_content(input, stream_url.as_str(), stream_file_path) + ) { + (Ok(category_content), Ok(stream_content)) => { + match xtream_parser::parse_xtream(input, + *xtream_cluster, + &category_content, + &stream_content) { + Ok(sub_playlist_parsed) => { + if let Some(mut xtream_sub_playlist) = sub_playlist_parsed { + playlist_groups.extend(xtream_sub_playlist.drain(..)); } } Err(err) => errors.push(err) } - } - Err(err) => errors.push(err) + }, + (Err(err1), Err(err2)) => { + errors.extend([err1, err2]); + }, + (Err(err), _) => errors.push(err), + (_, Err(err)) => errors.push(err), } } } - playlist.sort_by(|a, b| a.title.partial_cmp(&b.title).unwrap_or(Ordering::Greater)); + playlist_groups.sort_by(|a, b| a.title.partial_cmp(&b.title).unwrap_or(Ordering::Greater)); - for (grp_id, plg) in (1_u32..).zip(playlist.iter_mut()) { + for (grp_id, plg) in (1_u32..).zip(playlist_groups.iter_mut()) { plg.id = grp_id; } - (playlist, errors) + (playlist_groups, errors) } pub(crate) async fn get_xmltv(_cfg: &Config, input: &ConfigInput, working_dir: &String) -> (Option, Vec) { diff --git a/src/utils/request_utils.rs b/src/utils/request_utils.rs index 10b282c86..9d13c5c1a 100644 --- a/src/utils/request_utils.rs +++ b/src/utils/request_utils.rs @@ -232,15 +232,15 @@ pub(crate) async fn get_input_json_content(input: &ConfigInput, url: &str, persi Err(e) => create_m3u_filter_error_result!(M3uFilterErrorKind::Notify, "cant download input url: {url} => {}", e) } } - -pub(crate) fn get_base_url(url: &str) -> Option { - if let Some((scheme_end, rest)) = url.split_once("://") { - let scheme = scheme_end; - if let Some(authority_end) = rest.find('/') { - let authority = &rest[..authority_end]; - return Some(format!("{}://{}", scheme, authority)); - } - return Some(format!("{}://{}", scheme, rest)); - } - None -} +// +// pub(crate) fn get_base_url(url: &str) -> Option { +// if let Some((scheme_end, rest)) = url.split_once("://") { +// let scheme = scheme_end; +// if let Some(authority_end) = rest.find('/') { +// let authority = &rest[..authority_end]; +// return Some(format!("{}://{}", scheme, authority)); +// } +// return Some(format!("{}://{}", scheme, rest)); +// } +// None +// }