diff --git a/Cargo.lock b/Cargo.lock index 9b2071734..b0d581d34 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -87,7 +87,7 @@ dependencies = [ "mime", "percent-encoding", "pin-project-lite", - "rand", + "rand 0.8.5", "sha1", "smallvec", "tokio", @@ -264,7 +264,7 @@ dependencies = [ "getrandom 0.2.15", "once_cell", "version_check", - "zerocopy", + "zerocopy 0.7.35", ] [[package]] @@ -375,112 +375,10 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "81953c529336010edd6d8e358f886d9581267795c61b19475b71314bffa46d35" dependencies = [ "concurrent-queue", - "event-listener 2.5.3", + "event-listener", "futures-core", ] -[[package]] -name = "async-channel" -version = "2.3.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "89b47800b0be77592da0afd425cc03468052844aff33b84e33cc696f64e77b6a" -dependencies = [ - "concurrent-queue", - "event-listener-strategy", - "futures-core", - "pin-project-lite", -] - -[[package]] -name = "async-executor" -version = "1.13.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "30ca9a001c1e8ba5149f91a74362376cc6bc5b919d92d988668657bd570bdcec" -dependencies = [ - "async-task", - "concurrent-queue", - "fastrand 2.3.0", - "futures-lite 2.6.0", - "slab", -] - -[[package]] -name = "async-global-executor" -version = "2.4.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "05b1b633a2115cd122d73b955eadd9916c18c8f510ec9cd1686404c60ad1c29c" -dependencies = [ - "async-channel 2.3.1", - "async-executor", - "async-io", - "async-lock", - "blocking", - "futures-lite 2.6.0", - "once_cell", -] - -[[package]] -name = "async-io" -version = "2.4.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "43a2b323ccce0a1d90b449fd71f2a06ca7faa7c54c2751f06c9bd851fc061059" -dependencies = [ - "async-lock", - "cfg-if", - "concurrent-queue", - "futures-io", - "futures-lite 2.6.0", - "parking", - "polling 3.7.4", - "rustix", - "slab", - "tracing", - "windows-sys 0.59.0", -] - -[[package]] -name = "async-lock" -version = "3.4.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "ff6e472cdea888a4bd64f342f09b3f50e1886d32afe8df3d663c01140b811b18" -dependencies = [ - "event-listener 5.4.0", - "event-listener-strategy", - "pin-project-lite", -] - -[[package]] -name = "async-std" -version = "1.13.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "c634475f29802fde2b8f0b505b1bd00dfe4df7d4a000f0b36f7671197d5c3615" -dependencies = [ - "async-channel 1.9.0", - "async-global-executor", - "async-io", - "async-lock", - "crossbeam-utils", - "futures-channel", - "futures-core", - "futures-io", - "futures-lite 2.6.0", - "gloo-timers", - "kv-log-macro", - "log", - "memchr", - "once_cell", - "pin-project-lite", - "pin-utils", - "slab", - "wasm-bindgen-futures", -] - -[[package]] -name = "async-task" -version = "4.7.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "8b75356056920673b02621b35afd0f7dda9306d03c79a30f5c56c44cf256e3de" - [[package]] name = "atomic-waker" version = "1.1.2" @@ -543,9 +441,9 @@ checksum = "8f68f53c83ab957f72c32642f3868eec03eb974d1fb82e453128456482613d36" [[package]] name = "blake2b_simd" -version = "1.0.2" +version = "1.0.3" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "23285ad32269793932e830392f2fe2f83e26488fd3ec778883a93c8323735780" +checksum = "06e903a20b159e944f91ec8499fe1e55651480c541ea0a584f5d967c49ad9d99" dependencies = [ "arrayref", "arrayvec", @@ -574,19 +472,6 @@ dependencies = [ "generic-array", ] -[[package]] -name = "blocking" -version = "1.6.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "703f41c54fc768e63e091340b424302bb1c29ef4aa0c7f10fe849dfb114d29ea" -dependencies = [ - "async-channel 2.3.1", - "async-task", - "futures-io", - "futures-lite 2.6.0", - "piper", -] - [[package]] name = "brotli" version = "6.0.0" @@ -643,9 +528,9 @@ checksum = "a2698f953def977c68f935bb0dfa959375ad4638570e969e2f1e9f433cbf1af6" [[package]] name = "cc" -version = "1.2.11" +version = "1.2.12" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e4730490333d58093109dc02c23174c3f4d490998c3fed3cc8e82d57afedb9cf" +checksum = "755717a7de9ec452bf7f3f1a3099085deabd7f2962b861dae91ecd7a365903d2" dependencies = [ "jobserver", "libc", @@ -965,27 +850,6 @@ version = "2.5.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "0206175f82b8d6bf6652ff7d71a1e27fd2e4efde587fd368662814d6ec1d9ce0" -[[package]] -name = "event-listener" -version = "5.4.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "3492acde4c3fc54c845eaab3eed8bd00c7a7d881f78bfc801e43a93dec1331ae" -dependencies = [ - "concurrent-queue", - "parking", - "pin-project-lite", -] - -[[package]] -name = "event-listener-strategy" -version = "0.5.3" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "3c3e4e0dd3673c1139bf041f3008816d9cf2946bbfac2945c09e523b8d7b05b2" -dependencies = [ - "event-listener 5.4.0", - "pin-project-lite", -] - [[package]] name = "fastrand" version = "1.9.0" @@ -1116,19 +980,6 @@ dependencies = [ "waker-fn", ] -[[package]] -name = "futures-lite" -version = "2.6.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "f5edaec856126859abb19ed65f39e90fea3a9574b9707f13539acf4abf7eb532" -dependencies = [ - "fastrand 2.3.0", - "futures-core", - "futures-io", - "parking", - "pin-project-lite", -] - [[package]] name = "futures-macro" version = "0.3.31" @@ -1211,18 +1062,6 @@ version = "0.31.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "07e28edb80900c19c28f1072f2e8aeca7fa06b23cd4169cefe1af5aa3260783f" -[[package]] -name = "gloo-timers" -version = "0.3.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "bbb143cf96099802033e0d4f4963b19fd2e0b728bcf076cd9cf7f6634f092994" -dependencies = [ - "futures-channel", - "futures-core", - "js-sys", - "wasm-bindgen", -] - [[package]] name = "h2" version = "0.3.26" @@ -1273,12 +1112,6 @@ version = "0.5.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "2304e00983f87ffb38b55b444b5e3b60a884b5d30c0fca7d82fe33449bbe55ea" -[[package]] -name = "hermit-abi" -version = "0.4.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "fbf6a919d6cf397374f7dfeeea91d974c7c0a7221d0d0f4f20d859d329e53fcc" - [[package]] name = "http" version = "0.2.12" @@ -1626,19 +1459,19 @@ version = "1.7.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "334e04b4d781f436dc315cb1e7515bd96826426345d498149e4bde36b67f8ee9" dependencies = [ - "async-channel 1.9.0", + "async-channel", "castaway", "crossbeam-utils", "curl", "curl-sys", "encoding_rs", - "event-listener 2.5.3", - "futures-lite 1.13.0", + "event-listener", + "futures-lite", "http 0.2.12", "log", "mime", "once_cell", - "polling 2.8.0", + "polling", "serde", "serde_json", "slab", @@ -1689,15 +1522,6 @@ dependencies = [ "simple_asn1", ] -[[package]] -name = "kv-log-macro" -version = "1.0.7" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "0de8b303297635ad57c9f5059fd9cee7a47f8e8daa09df0fcd07dd39fb22977f" -dependencies = [ - "log", -] - [[package]] name = "language-tags" version = "0.3.2" @@ -1787,9 +1611,6 @@ name = "log" version = "0.4.25" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "04cbf5b083de1c7e0222a7a51dbfdba1cbe1c6ab0b15e29fff3f6c077fd9cd9f" -dependencies = [ - "value-bag", -] [[package]] name = "m3u-filter" @@ -1798,10 +1619,8 @@ dependencies = [ "actix-cors", "actix-files", "actix-rt", - "actix-server", "actix-web", "actix-web-httpauth", - "async-std", "bincode", "blake3", "bytes", @@ -1817,13 +1636,13 @@ dependencies = [ "libc", "log", "mime", - "openssl", + "parking_lot", "paste", "path-clean", "pest", "pest_derive", "quick-xml", - "rand", + "rand 0.9.0", "regex", "reqwest", "rpassword", @@ -1834,7 +1653,6 @@ dependencies = [ "serde_json", "serde_yaml", "tempfile", - "time", "tokio", "tokio-stream", "unidecode", @@ -1983,15 +1801,6 @@ version = "0.1.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d05e27ee213611ffe7d6348b942e8f942b37114c00cc03cec254295a4a17852e" -[[package]] -name = "openssl-src" -version = "300.4.1+3.4.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "faa4eac4138c62414b5622d1b31c5c304f34b406b013c079c2bbc652fdd6678c" -dependencies = [ - "cc", -] - [[package]] name = "openssl-sys" version = "0.9.105" @@ -2000,7 +1809,6 @@ checksum = "8b22d5b84be05a8d6947c7cb71f7c849aa0f112acd4bf51c2a7c1c988ac0a9dc" dependencies = [ "cc", "libc", - "openssl-src", "pkg-config", "vcpkg", ] @@ -2139,17 +1947,6 @@ version = "0.1.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "8b870d8c151b6f2fb93e84a13146138f05d02ed11c7e7c54f8826aaaf7c9f184" -[[package]] -name = "piper" -version = "0.2.4" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "96c8c490f422ef9a4efd2cb5b42b76c8613d7e7dfc1caf667b8a3350a5acc066" -dependencies = [ - "atomic-waker", - "fastrand 2.3.0", - "futures-io", -] - [[package]] name = "pkg-config" version = "0.3.31" @@ -2172,21 +1969,6 @@ dependencies = [ "windows-sys 0.48.0", ] -[[package]] -name = "polling" -version = "3.7.4" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "a604568c3202727d1507653cb121dbd627a58684eb09a820fd746bee38b4442f" -dependencies = [ - "cfg-if", - "concurrent-queue", - "hermit-abi", - "pin-project-lite", - "rustix", - "tracing", - "windows-sys 0.59.0", -] - [[package]] name = "powerfmt" version = "0.2.0" @@ -2199,7 +1981,7 @@ version = "0.2.20" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "77957b295656769bb8ad2b6a6b09d897d94f05c41b069aede1fcdaa675eaea04" dependencies = [ - "zerocopy", + "zerocopy 0.7.35", ] [[package]] @@ -2247,7 +2029,7 @@ checksum = "a2fe5ef3495d7d2e377ff17b1a8ce2ee2ec2a18cde8b6ad6619d65d0701c135d" dependencies = [ "bytes", "getrandom 0.2.15", - "rand", + "rand 0.8.5", "ring", "rustc-hash", "rustls", @@ -2289,8 +2071,19 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "34af8d1a0e25924bc5b7c43c079c942339d8f0a8b57c39049bef581b46327404" dependencies = [ "libc", - "rand_chacha", - "rand_core", + "rand_chacha 0.3.1", + "rand_core 0.6.4", +] + +[[package]] +name = "rand" +version = "0.9.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3779b94aeb87e8bd4e834cee3650289ee9e0d5677f976ecdb6d219e5f4f6cd94" +dependencies = [ + "rand_chacha 0.9.0", + "rand_core 0.9.0", + "zerocopy 0.8.16", ] [[package]] @@ -2300,7 +2093,17 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "e6c10a63a0fa32252be49d21e7709d4d4baf8d231c2dbce1eaa8141b9b127d88" dependencies = [ "ppv-lite86", - "rand_core", + "rand_core 0.6.4", +] + +[[package]] +name = "rand_chacha" +version = "0.9.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d3022b5f1df60f26e1ffddd6c66e8aa15de382ae63b3a0c1bfc0e4d3e3f325cb" +dependencies = [ + "ppv-lite86", + "rand_core 0.9.0", ] [[package]] @@ -2312,6 +2115,16 @@ dependencies = [ "getrandom 0.2.15", ] +[[package]] +name = "rand_core" +version = "0.9.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b08f3c9802962f7e1b25113931d94f43ed9725bebc59db9d0c3e9a23b67e15ff" +dependencies = [ + "getrandom 0.3.1", + "zerocopy 0.8.16", +] + [[package]] name = "redox_syscall" version = "0.5.8" @@ -2463,9 +2276,9 @@ checksum = "719b953e2095829ee67db738b3bfa9fa368c94900df327b3f07fe6e794d2fe1f" [[package]] name = "rustc-hash" -version = "2.1.0" +version = "2.1.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "c7fb8039b3032c191086b10f11f319a6e99e1e82889c5cc6046f515c9db1d497" +checksum = "357703d41365b4b27c590e3ed91eabb1b663f07c4c084095e60cbed4362dff0d" [[package]] name = "rustc_version" @@ -2731,7 +2544,7 @@ version = "0.5.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "6d7400c0eff44aa2fcb5e31a5f24ba9716ed90138769e4977a2ba6014ae63eb5" dependencies = [ - "async-channel 1.9.0", + "async-channel", "futures-core", "futures-io", ] @@ -3138,12 +2951,6 @@ version = "0.15.8" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "4e8257fbc510f0a46eb602c10215901938b5c2a7d5e70fc11483b1d3c9b5b18c" -[[package]] -name = "value-bag" -version = "1.10.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "3ef4c4aa54d5d05a279399bfa921ec387b7aba77caf7a682ae8d86785b8fdad2" - [[package]] name = "vcpkg" version = "0.2.15" @@ -3569,7 +3376,16 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "1b9b4fd18abc82b8136838da5d50bae7bdea537c574d8dc1a34ed098d6c166f0" dependencies = [ "byteorder", - "zerocopy-derive", + "zerocopy-derive 0.7.35", +] + +[[package]] +name = "zerocopy" +version = "0.8.16" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7b8c07a70861ce02bad1607b5753ecb2501f67847b9f9ada7c160fff0ec6300c" +dependencies = [ + "zerocopy-derive 0.8.16", ] [[package]] @@ -3583,6 +3399,17 @@ dependencies = [ "syn", ] +[[package]] +name = "zerocopy-derive" +version = "0.8.16" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5226bc9a9a9836e7428936cde76bb6b22feea1a8bfdbc0d241136e4d13417e25" +dependencies = [ + "proc-macro2", + "quote", + "syn", +] + [[package]] name = "zerofrom" version = "0.1.5" diff --git a/Cargo.toml b/Cargo.toml index 86d67d2be..a0568cb3a 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -25,7 +25,6 @@ reqwest = { version = "0", features = ["blocking", "json", "stream", "rustls-tls chrono = "0.4" cron = "0.15" actix-web = "4.9" -actix-server = "2.5" actix-files = "0" actix-cors = "0" actix-rt = "2.10" @@ -38,25 +37,23 @@ pest = "2.7" pest_derive = "2.7" enum-iterator = "2" unidecode = "0" -openssl = { version = "*", features = ["vendored"] } #https://docs.rs/openssl/0.10.34/openssl/#vendored mime = "0.3" log = "0.4" env_logger = "0.11" rustelebot = "0.3" bincode = "1.3" -rand = "0.8" +rand = "0.9" rpassword = "7.3" flate2 = "1" -time = "0.3" blake3 = "1.5" -bytes = "1.9" -async-std = "1.13" +bytes = "1.10" tokio-stream = { version = "0.1", features = ["sync"] } tokio = "1.43" paste = "1.0" tempfile = "3.15" ruzstd = "0" filetime = "0.2" +parking_lot = "0.12" #[cfg(target_os = "macos")] libc = "0" #[cfg(target_os = "windows")] diff --git a/src/api/api_utils.rs b/src/api/api_utils.rs index 0463d6340..a019d7be4 100644 --- a/src/api/api_utils.rs +++ b/src/api/api_utils.rs @@ -18,7 +18,7 @@ use actix_files::NamedFile; use actix_web::body::{BodyStream, SizedStream}; use actix_web::http::header::{HeaderValue, CACHE_CONTROL}; use actix_web::{HttpRequest, HttpResponse}; -use async_std::sync::Mutex; +use parking_lot::FairMutex; use futures::{TryStreamExt}; use log::{error, log_enabled, trace}; use reqwest::StatusCode; @@ -86,13 +86,13 @@ pub async fn serve_file(file_path: &Path, req: &HttpRequest, mime_type: mime::Mi pub async fn get_user_target_by_credentials<'a>(username: &str, password: &str, api_req: &'a UserApiRequest, app_state: &'a AppState) -> Option<(ProxyUserCredentials, &'a ConfigTarget)> { if !username.is_empty() && !password.is_empty() { - app_state.config.get_target_for_user(username, password).await + app_state.config.get_target_for_user(username, password) } else { let token = api_req.token.as_str().trim(); if token.is_empty() { None } else { - app_state.config.get_target_for_user_by_token(token).await + app_state.config.get_target_for_user_by_token(token) } } } @@ -137,7 +137,7 @@ pub async fn stream_response(app_state: &AppState, stream_url: &str, let share_stream = is_stream_share_enabled(item_type, target); if share_stream { - if let Some(value) = shared_stream_response(app_state, stream_url).await { + if let Some(value) = shared_stream_response(app_state, stream_url) { return value; } } @@ -150,7 +150,7 @@ pub async fn stream_response(app_state: &AppState, stream_url: &str, let (stream_opt, provider_response) = if direct_pipe_provider_stream { get_provider_pipe_stream(&app_state.http_client, &url, req, input).await } else { - let buffer_stream_options = BufferStreamOptions::new(item_type, stream_retry, buffer_enabled, buffer_size); + let buffer_stream_options = BufferStreamOptions::new(item_type, stream_retry, buffer_enabled, buffer_size, share_stream); provider_stream::get_provider_reconnect_buffered_stream(&app_state.http_client, &url, req, input, buffer_stream_options).await }; if let Some(stream) = stream_opt { @@ -159,10 +159,14 @@ pub async fn stream_response(app_state: &AppState, stream_url: &str, let stream = ActiveClientStream::new(stream, active_clients, log_active_clients); let stream_resp = if share_stream { let shared_headers = provider_response.as_ref().map_or_else(Vec::new, |(h, _)| h.clone()); - SharedStreamManager::subscribe(app_state, stream_url, stream, shared_stream_use_own_buffer, shared_headers).await; - if let Some(broadcast_stream) = SharedStreamManager::subscribe_shared_stream(app_state, stream_url).await { + SharedStreamManager::subscribe(app_state, stream_url, stream, shared_stream_use_own_buffer, shared_headers); + if let Some(broadcast_stream) = SharedStreamManager::subscribe_shared_stream(app_state, stream_url) { let mut response_builder = get_stream_response_with_headers(provider_response, stream_url); - if content_length > 0 { response_builder.body(SizedStream::new(content_length, broadcast_stream)) } else { response_builder.body(BodyStream::new(broadcast_stream)) } + if content_length > 0 { + response_builder.body(SizedStream::new(content_length, broadcast_stream)) } + else { + response_builder.body(BodyStream::new(broadcast_stream)) + } } else { HttpResponse::BadRequest().finish() } @@ -178,10 +182,10 @@ pub async fn stream_response(app_state: &AppState, stream_url: &str, HttpResponse::BadRequest().finish() } -async fn shared_stream_response(app_state: &AppState, stream_url: &str) -> Option { - if let Some(stream) = SharedStreamManager::subscribe_shared_stream(app_state, stream_url).await { +fn shared_stream_response(app_state: &AppState, stream_url: &str) -> Option { + if let Some(stream) = SharedStreamManager::subscribe_shared_stream(app_state, stream_url) { debug_if_enabled!("Using shared channel {}", sanitize_sensitive_info(stream_url)); - if let Some(headers) = app_state.shared_stream_manager.lock().await.get_shared_state_headers(stream_url).await { + if let Some(headers) = app_state.shared_stream_manager.lock().get_shared_state_headers(stream_url) { let mut response_builder = get_stream_response_with_headers(Some((headers.clone(), StatusCode::OK)), stream_url); return Some(response_builder.body(BodyStream::new(stream))); } @@ -205,7 +209,7 @@ pub fn get_headers_from_request(req: &HttpRequest, filter: &HeaderFilter) -> Has .collect() } -fn get_add_cache_content(res_url: &str, cache: &Arc>>) -> Box { +fn get_add_cache_content(res_url: &str, cache: &Arc>>) -> Box { let resource_url = String::from(res_url); let cache = Arc::clone(cache); let add_cache_content: Box = Box::new(move |size| { @@ -213,8 +217,8 @@ fn get_add_cache_content(res_url: &str, cache: &Arc Option { @@ -47,7 +47,7 @@ async fn save_config_api_proxy_user( ) -> HttpResponse { let mut users = req.0; users.iter_mut().flat_map(|t| &mut t.credentials).for_each(ProxyUserCredentials::trim); - if let Some(api_proxy) = app_state.config.t_api_proxy.write().await.as_mut() { + if let Some(api_proxy) = app_state.config.t_api_proxy.write().as_mut() { let backup_dir = app_state.config.backup_dir.as_ref().unwrap().as_str(); api_proxy.user = users; if let Some(err) = intern_save_config_api_proxy(backup_dir, api_proxy, app_state.config.t_api_proxy_file_path.as_str()) { @@ -85,7 +85,7 @@ async fn save_config_api_proxy_config( return HttpResponse::BadRequest().json(json!({"error": "Invalid content"})); } } - if let Some(api_proxy) = app_state.config.t_api_proxy.write().await.as_mut() { + if let Some(api_proxy) = app_state.config.t_api_proxy.write().as_mut() { api_proxy.server = req_api_proxy; let backup_dir = app_state.config.backup_dir.as_ref().unwrap().as_str(); if let Some(err) = intern_save_config_api_proxy(backup_dir, api_proxy, app_state.config.t_api_proxy_file_path.as_str()) { @@ -227,7 +227,7 @@ async fn config( // if we didn't read it from file then we should use it from app_state if result.api_proxy.is_none() { - result.api_proxy.clone_from(&*app_state.config.t_api_proxy.read().await); + result.api_proxy.clone_from(&*app_state.config.t_api_proxy.read()); } HttpResponse::Ok().json(result) diff --git a/src/api/endpoints/xtream_api.rs b/src/api/endpoints/xtream_api.rs index d1ae2a501..0daef24e7 100644 --- a/src/api/endpoints/xtream_api.rs +++ b/src/api/endpoints/xtream_api.rs @@ -132,8 +132,8 @@ fn get_xtream_player_api_stream_url(input: &ConfigInput, context: &XtreamApiStre } } -async fn get_user_info(user: &ProxyUserCredentials, cfg: &Config) -> XtreamAuthorizationResponse { - let server_info = cfg.get_user_server_info(user).await; +fn get_user_info(user: &ProxyUserCredentials, cfg: &Config) -> XtreamAuthorizationResponse { + let server_info = cfg.get_user_server_info(user); XtreamAuthorizationResponse::new(&server_info, user) } @@ -151,7 +151,7 @@ async fn xtream_player_api_stream( } let (action_stream_id, stream_ext) = separate_number_and_remainder(stream_req.stream_id); let virtual_id: u32 = try_result_bad_request!(action_stream_id.trim().parse()); - let (pli, mapping) = try_result_bad_request!(xtream_repository::xtream_get_item_for_stream_id(virtual_id, &app_state.config, target, None).await, true, format!("Failed to read xtream item for stream id {}", virtual_id)); + let (pli, mapping) = try_result_bad_request!(xtream_repository::xtream_get_item_for_stream_id(virtual_id, &app_state.config, target, None), true, format!("Failed to read xtream item for stream id {}", virtual_id)); let input = try_option_bad_request!(app_state.config.get_input_by_name(pli.input_name.as_str()), true, format!("Cant find input for target {target_name}, context {}, stream_id {virtual_id}", stream_req.context)); if pli.item_type == PlaylistItemType::LiveHls { @@ -222,13 +222,13 @@ fn get_doc_resource_field_value(field: &str, doc: Option<&Value>) -> Option Result>, serde_json::Error> { +fn xtream_get_info_resource_url(config: &Config, pli: &XtreamPlaylistItem, target: &ConfigTarget, resource: &str) -> Result>, serde_json::Error> { let info_content = match pli.xtream_cluster { XtreamCluster::Video => { - xtream_repository::xtream_load_vod_info(config, target.name.as_str(), pli.get_virtual_id()).await + xtream_repository::xtream_load_vod_info(config, target.name.as_str(), pli.get_virtual_id()) } XtreamCluster::Series => { - xtream_repository::xtream_load_series_info(config, target.name.as_str(), pli.get_virtual_id()).await + xtream_repository::xtream_load_series_info(config, target.name.as_str(), pli.get_virtual_id()) } XtreamCluster::Live => None, }; @@ -297,10 +297,10 @@ fn get_season_info_doc(doc: &Vec, season_id: u32) -> Option<&Value> { } -async fn xtream_get_season_resource_url(config: &Config, pli: &XtreamPlaylistItem, target: &ConfigTarget, resource: &str) -> Result>, serde_json::Error> { +fn xtream_get_season_resource_url(config: &Config, pli: &XtreamPlaylistItem, target: &ConfigTarget, resource: &str) -> Result>, serde_json::Error> { let info_content = match pli.xtream_cluster { XtreamCluster::Series => { - xtream_repository::xtream_load_series_info(config, target.name.as_str(), pli.get_virtual_id()).await + xtream_repository::xtream_load_series_info(config, target.name.as_str(), pli.get_virtual_id()) } XtreamCluster::Video | XtreamCluster::Live => None, }; @@ -339,11 +339,11 @@ async fn xtream_player_api_resource( } let virtual_id: u32 = try_result_bad_request!(resource_req.stream_id.trim().parse()); let resource = resource_req.action_path.trim(); - let (pli, _) = try_result_bad_request!(xtream_repository::xtream_get_item_for_stream_id(virtual_id, &app_state.config, target, None).await, true, format!("Failed to read xtream item for stream id {}", virtual_id)); + let (pli, _) = try_result_bad_request!(xtream_repository::xtream_get_item_for_stream_id(virtual_id, &app_state.config, target, None), true, format!("Failed to read xtream item for stream id {}", virtual_id)); let stream_url = if resource.starts_with(INFO_RESOURCE_PREFIX) { - try_result_bad_request!(xtream_get_info_resource_url(&app_state.config, &pli, target, resource).await) + try_result_bad_request!(xtream_get_info_resource_url(&app_state.config, &pli, target, resource)) } else if resource.starts_with(SEASON_RESOURCE_PREFIX) { - try_result_bad_request!(xtream_get_season_resource_url(&app_state.config, &pli, target, resource).await) + try_result_bad_request!(xtream_get_season_resource_url(&app_state.config, &pli, target, resource)) } else { pli.get_field(resource) }; @@ -434,7 +434,7 @@ async fn xtream_get_stream_info_response(app_state: &AppState, user: &ProxyUserC Err(_) => return HttpResponse::BadRequest().finish() }; - if let Ok((pli, virtual_record)) = xtream_repository::xtream_get_item_for_stream_id(virtual_id, &app_state.config, target, Some(cluster)).await { + if let Ok((pli, virtual_record)) = xtream_repository::xtream_get_item_for_stream_id(virtual_id, &app_state.config, target, Some(cluster)) { if pli.provider_id > 0 { let input_name = Rc::clone(&pli.input_name); if let Some(input) = app_state.config.get_input_by_name(input_name.as_str()) { @@ -468,7 +468,7 @@ async fn xtream_get_short_epg(app_state: &AppState, user: &ProxyUserCredentials, Err(_) => return HttpResponse::BadRequest().finish() }; - if let Ok((pli, _)) = xtream_repository::xtream_get_item_for_stream_id(virtual_id, &app_state.config, target, None).await { + if let Ok((pli, _)) = xtream_repository::xtream_get_item_for_stream_id(virtual_id, &app_state.config, target, None) { if pli.provider_id > 0 { let input_name = Rc::clone(&pli.input_name); if let Some(input) = app_state.config.get_input_by_name(input_name.as_str()) { @@ -520,7 +520,7 @@ async fn xtream_player_api_handle_content_action(config: &Config, target_name: & async fn xtream_get_catchup_response(app_state: &AppState, target: &ConfigTarget, stream_id: &str, start: &str, end: &str) -> HttpResponse { let virtual_id: u32 = try_result_bad_request!(FromStr::from_str(stream_id)); - let (pli, _) = try_result_bad_request!(xtream_repository::xtream_get_item_for_stream_id(virtual_id, &app_state.config, target, Some(XtreamCluster::Live)).await); + let (pli, _) = try_result_bad_request!(xtream_repository::xtream_get_item_for_stream_id(virtual_id, &app_state.config, target, Some(XtreamCluster::Live))); let input = try_option_bad_request!(app_state.config.get_input_by_name(pli.input_name.as_str())); let info_url = try_option_bad_request!(xtream::get_xtream_player_api_action_url(input, ACTION_GET_CATCHUP_TABLE).map(|action_url| format!("{action_url}&{TAG_STREAM_ID}={}&start={start}&end={end}", pli.provider_id))); let content = try_result_bad_request!(xtream::get_xtream_stream_info_content(Arc::clone(&app_state.http_client), info_url.as_str(), input).await); @@ -528,13 +528,7 @@ async fn xtream_get_catchup_response(app_state: &AppState, target: &ConfigTarget let epg_listings = try_option_bad_request!(doc.get_mut(TAG_EPG_LISTINGS).and_then(Value::as_array_mut)); let target_path = try_option_bad_request!(get_target_storage_path(&app_state.config, target.name.as_str())); let mut target_id_mapping = { - let _file_lock = match app_state.config.file_locks.read_lock(&target_path).await { - Ok(lock) => lock, - Err(err) => { - error!("Could not get lock for id mapping for target {} err:{err}", target.name); - return HttpResponse::InternalServerError().finish(); - } - }; + let _file_lock = app_state.config.file_locks.read_lock(&target_path); TargetIdMapping::new(&target_path) }; for epg_list_item in epg_listings.iter_mut().filter_map(Value::as_object_mut) { @@ -546,13 +540,7 @@ async fn xtream_get_catchup_response(app_state: &AppState, target: &ConfigTarget } } { - let _file_lock = match app_state.config.file_locks.write_lock(&target_path).await { - Ok(lock) => lock, - Err(err) => { - error!("Could not get lock for id mapping for target {} err:{err}", target.name); - return HttpResponse::InternalServerError().finish(); - } - }; + let _file_lock = app_state.config.file_locks.write_lock(&target_path); if let Err(err) = target_id_mapping.persist() { error!("Failed to write catchup id mapping {err}"); return HttpResponse::BadRequest().finish(); @@ -588,12 +576,12 @@ async fn xtream_player_api( let user_target = get_user_target(&api_req, app_state).await; if let Some((user, target)) = user_target { if !target.has_output(&TargetType::Xtream) { - return HttpResponse::Ok().json(get_user_info(&user, &app_state.config).await); + return HttpResponse::Ok().json(get_user_info(&user, &app_state.config)); } let action = api_req.action.trim(); if action.is_empty() { - return HttpResponse::Ok().json(get_user_info(&user, &app_state.config).await); + return HttpResponse::Ok().json(get_user_info(&user, &app_state.config)); } // Process specific playlist actions @@ -634,11 +622,11 @@ async fn xtream_player_api( let category_id = api_req.category_id.trim().parse::().unwrap_or(0); let result = match action { ACTION_GET_LIVE_STREAMS => - skip_flag_optional!(skip_live, xtream_repository::xtream_load_rewrite_playlist(XtreamCluster::Live, &app_state.config, target, category_id, &user).await), + skip_flag_optional!(skip_live, xtream_repository::xtream_load_rewrite_playlist(XtreamCluster::Live, &app_state.config, target, category_id, &user)), ACTION_GET_VOD_STREAMS => - skip_flag_optional!(skip_vod, xtream_repository::xtream_load_rewrite_playlist(XtreamCluster::Video, &app_state.config, target, category_id, &user).await), + skip_flag_optional!(skip_vod, xtream_repository::xtream_load_rewrite_playlist(XtreamCluster::Video, &app_state.config, target, category_id, &user)), ACTION_GET_SERIES => - skip_flag_optional!(skip_series, xtream_repository::xtream_load_rewrite_playlist(XtreamCluster::Series, &app_state.config, target, category_id, &user).await), + skip_flag_optional!(skip_series, xtream_repository::xtream_load_rewrite_playlist(XtreamCluster::Series, &app_state.config, target, category_id, &user)), _ => Some(Err(info_err!(format!("Cant find action: {action} for target: {}", &target.name)) )), }; diff --git a/src/api/main_api.rs b/src/api/main_api.rs index 454418114..94a4c1f15 100644 --- a/src/api/main_api.rs +++ b/src/api/main_api.rs @@ -2,11 +2,12 @@ use actix_cors::Cors; use actix_web::middleware::Logger; use actix_web::web::Data; use actix_web::{web, App, HttpResponse, HttpServer}; -use async_std::sync::{Mutex, RwLock}; +use parking_lot::{FairMutex}; +use tokio::sync::{RwLock, Mutex}; use log::{error, info}; use std::collections::{VecDeque}; use std::io::ErrorKind; -use std::path::{PathBuf}; +use std::path::PathBuf; use std::sync::Arc; use std::sync::atomic::{AtomicUsize, Ordering}; use crate::api::endpoints::hls_api::hls_api_register; @@ -50,14 +51,14 @@ async fn healthcheck(app_state: web::Data,) -> HttpResponse { fn create_shared_data(cfg: &Arc) -> Data { let lru_cache = cfg.reverse_proxy.as_ref().and_then(|r| r.cache.as_ref()).and_then(|c| if c.enabled { - Some(Mutex::new(LRUResourceCache::new(c.t_size, &PathBuf::from(c.dir.as_ref().unwrap())))) + Some(FairMutex::new(LRUResourceCache::new(c.t_size, &PathBuf::from(c.dir.as_ref().unwrap())))) } else { None} ); let cache = Arc::new(lru_cache); let cache_scanner = Arc::clone(&cache); actix_rt::spawn(async move { if let Some(m) = cache_scanner.as_ref() { - let mut c = m.lock().await; - if let Err(err) = (*c).scan().await { + let mut c = m.lock(); + if let Err(err) = (*c).scan() { error!("Failed to scan cache {err}"); } } @@ -70,7 +71,7 @@ fn create_shared_data(cfg: &Arc) -> Data { finished: Arc::from(RwLock::new(Vec::new())), }), active_clients: Arc::new(AtomicUsize::new(0)), - shared_stream_manager: Arc::new(Mutex::new(SharedStreamManager::new())), + shared_stream_manager: Arc::new(FairMutex::new(SharedStreamManager::new())), http_client: Arc::new(reqwest::Client::new()), cache, }) diff --git a/src/api/model/app_state.rs b/src/api/model/app_state.rs index b05941e76..9cadd91f1 100644 --- a/src/api/model/app_state.rs +++ b/src/api/model/app_state.rs @@ -1,6 +1,6 @@ use std::sync::Arc; use std::sync::atomic::AtomicUsize; -use async_std::sync::{Mutex}; +use parking_lot::FairMutex; use crate::api::model::download::DownloadQueue; use crate::api::model::streams::shared_stream_manager::SharedStreamManager; use crate::model::config::{Config}; @@ -9,8 +9,8 @@ use crate::tools::lru_cache::LRUResourceCache; pub struct AppState { pub config: Arc, pub downloads: Arc, - pub shared_stream_manager: Arc>, + pub shared_stream_manager: Arc>, pub active_clients: Arc, pub http_client: Arc, - pub cache: Arc>> + pub cache: Arc>> } diff --git a/src/api/model/download.rs b/src/api/model/download.rs index c6bdd1087..0ae732645 100644 --- a/src/api/model/download.rs +++ b/src/api/model/download.rs @@ -1,8 +1,8 @@ use std::collections::VecDeque; use std::ffi::OsStr; use std::path::{Path, PathBuf}; -use std::sync::{Arc}; -use async_std::sync::{RwLock, Mutex}; +use tokio::sync::{RwLock, Mutex}; +use std::sync::Arc; use actix_web::web; use serde::{Deserialize, Serialize}; use unidecode::unidecode; diff --git a/src/api/model/streams/active_client_stream.rs b/src/api/model/streams/active_client_stream.rs index db167377f..67a658035 100644 --- a/src/api/model/streams/active_client_stream.rs +++ b/src/api/model/streams/active_client_stream.rs @@ -2,7 +2,7 @@ use crate::api::model::streams::provider_stream_factory::ResponseStream; use bytes::Bytes; use std::pin::Pin; use std::sync::atomic::{AtomicUsize, Ordering}; -use std::sync::{Arc}; +use std::sync::Arc; use std::task::{Poll}; use futures::{Stream}; use log::info; diff --git a/src/api/model/streams/client_stream.rs b/src/api/model/streams/client_stream.rs index d883096d3..eec269dae 100644 --- a/src/api/model/streams/client_stream.rs +++ b/src/api/model/streams/client_stream.rs @@ -2,7 +2,7 @@ use crate::api::model::streams::provider_stream_factory::ResponseStream; use bytes::Bytes; use std::pin::Pin; use std::sync::atomic::{AtomicUsize, Ordering}; -use std::sync::{Arc}; +use std::sync::Arc; use std::task::{Poll}; use futures::{Stream}; use crate::api::model::stream_error::StreamError; diff --git a/src/api/model/streams/provider_stream_factory.rs b/src/api/model/streams/provider_stream_factory.rs index e3434ad2a..e106289d8 100644 --- a/src/api/model/streams/provider_stream_factory.rs +++ b/src/api/model/streams/provider_stream_factory.rs @@ -12,7 +12,7 @@ use actix_web::HttpRequest; use bytes::Bytes; use futures::stream::{self, BoxStream}; use futures::{StreamExt, TryStreamExt}; -use log::warn; +use log::{error, warn}; use reqwest::header::{HeaderMap, RANGE}; use reqwest::StatusCode; use std::collections::HashMap; @@ -34,6 +34,7 @@ pub struct BufferStreamOptions { reconnect_enabled: bool, buffer_enabled: bool, buffer_size: usize, + share_stream: bool, } impl BufferStreamOptions { @@ -42,12 +43,14 @@ impl BufferStreamOptions { reconnect_enabled: bool, buffer_enabled: bool, buffer_size: usize, + share_stream: bool, ) -> Self { Self { item_type, reconnect_enabled, buffer_enabled, buffer_size, + share_stream } } @@ -170,6 +173,7 @@ fn get_client_stream_request_params( // we need the range bytes from client request for seek ing to the right position let req_range_start_bytes = get_request_range_start_bytes(&req_headers); req_headers.remove("range"); + // These are the configured headers for this input. let input_headers = input.map(|i| i.headers.clone()); // We merge configured input headers with the headers from the request. @@ -206,7 +210,10 @@ async fn provider_request(request_client: Arc, initial_info: bo } else { None }; - return Ok(Some((response.bytes_stream().map_err(|err| StreamError::reqwest(&err)).boxed(), response_info))); + return Ok(Some((response.bytes_stream().map_err(|err| { + error!("Failed to read response body: {err}"); + StreamError::reqwest(&err) + }).boxed(), response_info))); } Err(status) } @@ -229,7 +236,10 @@ async fn stream_provider(client: Arc, stream_options: ProviderS Ok(response) => { let status = response.status(); if status.is_success() { - return Some(response.bytes_stream().map_err(|err| StreamError::reqwest(&err)).boxed()); + return Some(response.bytes_stream().map_err(|err| { + error!("Stream error {err}"); + StreamError::reqwest(&err) + }).boxed()); } if status.is_client_error() { return None; @@ -314,31 +324,30 @@ pub async fn create_provider_stream(client: Arc, options: BufferStreamOptions) -> Option { let stream_options = create_provider_stream_options(stream_url, req, input, &options); - let client_stream_factory = |stream, reconnect, range_cnt| { + let client_stream_factory = |stream, reconnect_flag, range_cnt| { let stream = if stream_options.is_buffered() { BufferedStream::new(stream, stream_options.get_buffer_size(), stream_options.get_continue_flag_clone(), stream_url.as_str()).boxed() } else { stream }; - ClientStream::new(stream, reconnect, range_cnt, stream_options.get_url().as_str()).boxed() + ClientStream::new(stream, reconnect_flag, range_cnt, stream_options.get_url().as_str()).boxed() }; match get_initial_stream(Arc::clone(&client), &stream_options).await { Some((init_stream, info)) => { - let is_video_stream = if let Some((headers, _)) = &info { + let is_media_stream = if let Some((headers, _)) = &info { classify_content_type(headers) == MimeCategory::Video } else { true // don't know what it is but lets assume it is }; let continue_signal = stream_options.get_continue_flag_clone(); - if is_video_stream && stream_options.should_reconnect() { + if is_media_stream && stream_options.should_reconnect() { let client_signal = Arc::clone(&continue_signal); let stream_options_provider = stream_options.clone(); let unfold: ResponseStream = stream::unfold((), move |()| { let client = Arc::clone(&client); let stream_opts = stream_options_provider.clone(); - async move { let stream = stream_provider(client, stream_opts).await?; Some((stream, ())) @@ -378,7 +387,7 @@ mod tests { let url = url::Url::parse("https://info.cern.ch/hypertext/WWW/TheProject.html").unwrap(); let input = None; - let options = BufferStreamOptions::new(PlaylistItemType::Live, true, true, 0); + let options = BufferStreamOptions::new(PlaylistItemType::Live, true, true, 0, false); let value = create_provider_stream(Arc::clone(&client), &url, &req, input, options); let mut values = value.await; 'outer: while let Some((ref mut stream, info)) = values.as_mut() { diff --git a/src/api/model/streams/shared_stream_manager.rs b/src/api/model/streams/shared_stream_manager.rs index deed1464a..bab976c77 100644 --- a/src/api/model/streams/shared_stream_manager.rs +++ b/src/api/model/streams/shared_stream_manager.rs @@ -3,7 +3,7 @@ use crate::api::model::streams::provider_stream_factory::STREAM_QUEUE_SIZE; use crate::api::model::stream_error::StreamError; use crate::utils::debug_if_enabled; use crate::utils::network::request::sanitize_sensitive_info; -use async_std::sync::Mutex; +use parking_lot::FairMutex; use bytes::Bytes; use futures::stream::BoxStream; use futures::{Stream, StreamExt}; @@ -52,7 +52,7 @@ fn convert_stream(stream: BoxStream) -> BoxStream, buf_size: usize, - subscribers: Arc>>>, + subscribers: Arc>>>, } impl SharedStreamState { @@ -61,17 +61,17 @@ impl SharedStreamState { Self { headers, buf_size, - subscribers: Arc::new(Mutex::new(Vec::new())), + subscribers: Arc::new(FairMutex::new(Vec::new())), } } - async fn subscribe(&self) -> BoxStream<'static, Result> { + fn subscribe(&self) -> BoxStream<'static, Result> { let (tx, rx) = mpsc::channel(self.buf_size); - self.subscribers.lock().await.push(tx); + self.subscribers.lock().push(tx); convert_stream(ReceiverStream::new(rx).boxed()) } - fn broadcast(&self, stream_url: &str, bytes_stream: S, shared_streams: Arc>) + fn broadcast(&self, stream_url: &str, bytes_stream: S, shared_streams: Arc>) where S: Stream> + Unpin + 'static, { @@ -82,7 +82,7 @@ impl SharedStreamState { actix_rt::spawn(async move { while let Some(item) = source_stream.next().await { if let Ok(data) = item { - let mut subs = subscriber.lock().await; + let mut subs = subscriber.lock(); if subs.len() > 0 { (*subs).retain(|sender| { match sender.try_send(data.clone()) { @@ -94,16 +94,17 @@ impl SharedStreamState { } else { debug_if_enabled!("No active subscribers. Closing shared provider stream {}", sanitize_sensitive_info(&streaming_url)); // Cleanup for removing unused shared streams - shared_streams.lock().await.unregister(&streaming_url).await; + shared_streams.lock().unregister(&streaming_url); return; } } } + shared_streams.lock().unregister(&streaming_url); }); } } -type SharedStreamRegister = Arc>>; +type SharedStreamRegister = Arc>>; pub struct SharedStreamManager { shared_streams: SharedStreamRegister, @@ -112,19 +113,19 @@ pub struct SharedStreamManager { impl SharedStreamManager { pub(crate) fn new() -> Self { Self { - shared_streams: Arc::new(Mutex::new(HashMap::new())), + shared_streams: Arc::new(FairMutex::new(HashMap::new())), } } - pub async fn get_shared_state_headers(&self, stream_url: &str) -> Option> { - self.shared_streams.lock().await.get(stream_url).map(|s| s.headers.clone()) + pub fn get_shared_state_headers(&self, stream_url: &str) -> Option> { + self.shared_streams.lock().get(stream_url).map(|s| s.headers.clone()) } - async fn unregister(&self, stream_url: &str) { - self.shared_streams.lock().await.remove(stream_url); + fn unregister(&self, stream_url: &str) { + self.shared_streams.lock().remove(stream_url); } - pub(crate) async fn subscribe( + pub(crate) fn subscribe( app_state: &AppState, stream_url: &str, bytes_stream: S, @@ -138,26 +139,26 @@ impl SharedStreamManager { shared_state.broadcast(stream_url, bytes_stream, Arc::clone(&app_state.shared_stream_manager)); app_state .shared_stream_manager - .lock().await + .lock() .shared_streams - .lock().await + .lock() .insert(stream_url.to_string(), shared_state); debug_if_enabled!("Created shared provider stream {}", sanitize_sensitive_info(stream_url)); } /// Creates a broadcast notify stream for the given URL if a shared stream exists. - pub async fn subscribe_shared_stream( + pub fn subscribe_shared_stream( app_state: &AppState, stream_url: &str, ) -> Option>> { if let Some(shared_stream) = app_state .shared_stream_manager - .lock().await + .lock() .shared_streams - .lock().await + .lock() .get(stream_url) { debug_if_enabled!("Responding existing shared client stream {}", sanitize_sensitive_info(stream_url)); - Some(shared_stream.subscribe().await) + Some(shared_stream.subscribe()) } else { None } diff --git a/src/auth/password.rs b/src/auth/password.rs index 2e0a65edf..f5b9e2d19 100644 --- a/src/auth/password.rs +++ b/src/auth/password.rs @@ -1,9 +1,9 @@ -use rand::{Rng, distributions::Alphanumeric, rngs::OsRng}; +use rand::{Rng}; +use rand::distr::Alphanumeric; use crate::m3u_filter_error::str_to_io_error; fn generate_salt(length: usize) -> String { - let rng = OsRng; - let salt: String = rng + let salt: String = rand::rng() .sample_iter(&Alphanumeric) .take(length) .map(char::from) diff --git a/src/main.rs b/src/main.rs index 8d3210126..4abffc107 100644 --- a/src/main.rs +++ b/src/main.rs @@ -21,7 +21,7 @@ use crate::auth::password::generate_password; use crate::model::config::{validate_targets, Config, HealthcheckConfig, LogLevelConfig, ProcessTargets}; use crate::model::healthcheck::Healthcheck; use crate::processing::processor::playlist; -use crate::utils::config_reader; +use utils::file::config_reader; use crate::utils::file::file_utils; use crate::utils::network::request::set_sanitize_sensitive_info; use clap::Parser; diff --git a/src/model/api_proxy.rs b/src/model/api_proxy.rs index 6e861b823..bcb53e452 100644 --- a/src/model/api_proxy.rs +++ b/src/model/api_proxy.rs @@ -5,7 +5,7 @@ use std::str::FromStr; use enum_iterator::Sequence; use log::debug; use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind, create_m3u_filter_error_result, info_err}; -use crate::utils::config_reader; +use crate::utils::file::config_reader; #[derive(Debug, Clone, serde::Serialize, serde::Deserialize, Sequence, PartialEq, Eq)] pub enum ProxyType { @@ -28,10 +28,7 @@ impl ProxyType { impl Display for ProxyType { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - write!( - f, - "{}", - match self { + write!(f, "{}", match self { Self::Reverse => Self::REVERSE, Self::Redirect => Self::REDIRECT, } diff --git a/src/model/config.rs b/src/model/config.rs index 95bdc6fa5..60e1e35c4 100644 --- a/src/model/config.rs +++ b/src/model/config.rs @@ -7,8 +7,8 @@ use std::fs::File; use std::io::BufRead; use std::path::PathBuf; use std::str::FromStr; -use std::sync::{Arc}; -use async_std::sync::RwLock; +use parking_lot::RwLock; +use std::sync::Arc; use crate::auth::user::UserCredential; use log::{debug, error, warn}; @@ -22,7 +22,7 @@ use crate::messaging::MsgKind; use crate::model::api_proxy::{ApiProxyConfig, ApiProxyServerInfo, ProxyUserCredentials}; use crate::model::mapping::Mapping; use crate::model::mapping::Mappings; -use crate::utils::config_reader; +use crate::utils::file::config_reader; use crate::utils::default_utils::{default_as_default, default_as_true, default_as_two_u16}; use crate::utils::file::file_lock_manager::FileLockManager; use crate::utils::file::file_utils; @@ -504,9 +504,9 @@ impl FromStr for InputType { type Err = M3uFilterError; fn from_str(s: &str) -> Result { - if s.eq("m3u") { + if s.eq(Self::M3U) { Ok(Self::M3u) - } else if s.eq("xtream") { + } else if s.eq(Self::XTREAM) { Ok(Self::Xtream) } else { create_m3u_filter_error_result!(M3uFilterErrorKind::Info, "Unknown InputType: {}", s) @@ -1072,16 +1072,16 @@ impl Config { None } - pub async fn get_target_for_user(&self, username: &str, password: &str) -> Option<(ProxyUserCredentials, &ConfigTarget)> { - self.t_api_proxy.read().await.as_ref().and_then(|api_proxy| self.intern_get_target_for_user(api_proxy.get_target_name(username, password))) + pub fn get_target_for_user(&self, username: &str, password: &str) -> Option<(ProxyUserCredentials, &ConfigTarget)> { + self.t_api_proxy.read().as_ref().and_then(|api_proxy| self.intern_get_target_for_user(api_proxy.get_target_name(username, password))) } - pub async fn get_target_for_user_by_token(&self, token: &str) -> Option<(ProxyUserCredentials, &ConfigTarget)> { - self.t_api_proxy.read().await.as_ref().and_then(|api_proxy| self.intern_get_target_for_user(api_proxy.get_target_name_by_token(token))) + pub fn get_target_for_user_by_token(&self, token: &str) -> Option<(ProxyUserCredentials, &ConfigTarget)> { + self.t_api_proxy.read().as_ref().and_then(|api_proxy| self.intern_get_target_for_user(api_proxy.get_target_name_by_token(token))) } - pub async fn get_user_credentials(&self, username: &str) -> Option { - self.t_api_proxy.read().await.as_ref().and_then(|api_proxy| api_proxy.get_user_credentials(username)) + pub fn get_user_credentials(&self, username: &str) -> Option { + self.t_api_proxy.read().as_ref().and_then(|api_proxy| api_proxy.get_user_credentials(username)) } pub fn get_input_by_name(&self, input_name: &str) -> Option<&ConfigInput> { @@ -1273,8 +1273,8 @@ impl Config { /// # Panics /// /// Will panic if default server invalid - pub async fn get_user_server_info(&self, user: &ProxyUserCredentials) -> ApiProxyServerInfo { - let server_info_list = self.t_api_proxy.read().await.as_ref().unwrap().server.clone(); + pub fn get_user_server_info(&self, user: &ProxyUserCredentials) -> ApiProxyServerInfo { + let server_info_list = self.t_api_proxy.read().as_ref().unwrap().server.clone(); let server_info_name = user.server.as_ref().map_or("default", |server_name| server_name.as_str()); server_info_list.iter().find(|c| c.name.eq(server_info_name)).map_or_else(|| server_info_list.first().unwrap().clone(), Clone::clone) } diff --git a/src/processing/processor/playlist.rs b/src/processing/processor/playlist.rs index c6de08dae..015e5beb6 100644 --- a/src/processing/processor/playlist.rs +++ b/src/processing/processor/playlist.rs @@ -5,7 +5,7 @@ use crate::model::config::ConfigRename; use crate::utils::network::epg; use crate::utils::network::m3u; use crate::utils::network::xtream; -use async_std::sync::Mutex; +use parking_lot::Mutex; use core::cmp::Ordering; use std::cell::RefCell; use std::collections::{HashMap, HashSet}; @@ -399,7 +399,7 @@ async fn process_sources(client: Arc, config: Arc, user for (index, _) in config.sources.iter().enumerate() { // We're using the file lock this way on purpose let source_lock_path = PathBuf::from(format!("source_{index}")); - let Ok(update_lock) = config.file_locks.try_write_lock(&source_lock_path).await else { + let Ok(update_lock) = config.file_locks.try_write_lock(&source_lock_path) else { warn!("The update operation for the source at index {index} was skipped because an update is already in progress."); continue; }; @@ -414,9 +414,9 @@ async fn process_sources(client: Arc, config: Arc, user let process = move || { System::new().block_on(async { let (input_stats, target_stats, mut res_errors) = process_source(Arc::clone(&http_client), cfg, index, usr_trgts).await; - shared_errors.lock().await.append(&mut res_errors); + shared_errors.lock().append(&mut res_errors); let process_stats = SourceStats::new(input_stats, target_stats); - shared_stats.lock().await.push(process_stats); + shared_stats.lock().push(process_stats); }); }; handles.push(thread::spawn(process)); @@ -425,9 +425,9 @@ async fn process_sources(client: Arc, config: Arc, user } } else { let (input_stats, target_stats, mut res_errors) = process_source(Arc::clone(&client), cfg, index, usr_trgts).await; - shared_errors.lock().await.append(&mut res_errors); + shared_errors.lock().append(&mut res_errors); let process_stats = SourceStats::new(input_stats, target_stats); - shared_stats.lock().await.push(process_stats); + shared_stats.lock().push(process_stats); } drop(update_lock); } diff --git a/src/processing/processor/xtream.rs b/src/processing/processor/xtream.rs index 44fb7b388..1e7cb518a 100644 --- a/src/processing/processor/xtream.rs +++ b/src/processing/processor/xtream.rs @@ -1,8 +1,8 @@ +use crate::m3u_filter_error::{info_err, notify_err}; use crate::m3u_filter_error::{str_to_io_error, to_io_error, M3uFilterError, M3uFilterErrorKind}; use crate::model::config::{Config, ConfigInput}; use crate::model::playlist::{FetchedPlaylist, PlaylistEntry, PlaylistItem, PlaylistItemType, XtreamCluster}; use crate::repository::storage::get_input_storage_path; -use crate::m3u_filter_error::{info_err, notify_err}; use serde::{Deserialize, Serialize}; use std::collections::HashMap; use std::fs::File; @@ -107,16 +107,14 @@ where } }; - match cfg.file_locks.read_lock(&file_path).await { - Ok(file_lock) => { - if let Ok(info_records) = BPlusTree::::load(&file_path) { - info_records.iter().for_each(|(provider_id, record)| { - processed_info_ids.insert(*provider_id, extract_ts(record)); - }); - } - drop(file_lock); + { + let file_lock = cfg.file_locks.read_lock(&file_path); + if let Ok(info_records) = BPlusTree::::load(&file_path) { + info_records.iter().for_each(|(provider_id, record)| { + processed_info_ids.insert(*provider_id, extract_ts(record)); + }); } - Err(err) => errors.push(info_err!(format!("{err}"))), + drop(file_lock); } processed_info_ids } diff --git a/src/processing/processor/xtream_series.rs b/src/processing/processor/xtream_series.rs index c7dff3222..1a997c51f 100644 --- a/src/processing/processor/xtream_series.rs +++ b/src/processing/processor/xtream_series.rs @@ -139,10 +139,7 @@ async fn process_series_info( return result; }; - let Ok(_file_lock) = cfg.file_locks.read_lock(&info_path).await else { - errors.push(notify_err!("Could not lock input info file for series".to_string())); - return result; - }; + let _file_lock = cfg.file_locks.read_lock(&info_path); // Contains the Series Info with episode listing let Ok(mut info_reader) = IndexedDocumentReader::::new(&info_path, &idx_path) else { return result; }; diff --git a/src/repository/kodi_repository.rs b/src/repository/kodi_repository.rs index a5d2736d3..e9c3eed6d 100644 --- a/src/repository/kodi_repository.rs +++ b/src/repository/kodi_repository.rs @@ -14,11 +14,10 @@ use crate::utils::file::file_lock_manager::FileReadGuard; use crate::utils::file::file_utils; use crate::utils::network::request::extract_extension_from_url; use crate::m3u_filter_error::{create_m3u_filter_error_result, info_err, notify_err}; -use async_std::fs::{create_dir_all, read_dir, remove_dir, remove_file, File}; -use async_std::io::{BufReadExt, BufReader, BufWriter, ReadExt, WriteExt}; +use tokio::fs::{create_dir_all, read_dir, remove_dir, remove_file, File}; +use tokio::io::{AsyncBufReadExt, AsyncReadExt, AsyncWriteExt, BufReader, BufWriter}; use chrono::Datelike; use filetime::{set_file_times, FileTime}; -use futures::StreamExt; use log::{debug, error}; use regex::Regex; use serde::Serialize; @@ -156,7 +155,7 @@ fn trim_whitespace(pattern: &Regex, input: &str) -> String { pattern.replace_all(input, " ").to_string() } -async fn kodi_style_rename( +fn kodi_style_rename( cfg: &Config, strm_item_info: &StrmItemInfo, style: &KodiStyle, @@ -191,7 +190,6 @@ async fn kodi_style_rename( input_tmdb_indexes, strm_item_info.item_type, ) - .await } _ => None, } { @@ -282,7 +280,7 @@ enum InputTmdbIndexValue { } type InputTmdbIndexMap = HashMap>; -async fn get_tmdb_value( +fn get_tmdb_value( cfg: &Config, provider_id: Option, input_name: &str, @@ -313,7 +311,8 @@ async fn get_tmdb_value( if let Ok(Some(tmdb_path)) = get_input_storage_path(input, &cfg.working_dir) .map(|storage_path| xtream_get_record_file_path(&storage_path, item_type)) { - if let Ok(file_lock) = cfg.file_locks.read_lock(&tmdb_path).await { + { + let file_lock = cfg.file_locks.read_lock(&tmdb_path); match item_type { PlaylistItemType::Series => { if let Ok(tree) = @@ -430,7 +429,7 @@ fn extract_item_info(pli: &PlaylistItem) -> StrmItemInfo { async fn prepare_strm_output_directory(path: &Path) -> Result<(), M3uFilterError> { // Ensure the directory exists - if let Err(e) = async_std::fs::create_dir_all(path).await { + if let Err(e) = tokio::fs::create_dir_all(path).await { error!("Failed to create directory {path:?}: {e}"); return create_m3u_filter_error_result!( M3uFilterErrorKind::Notify, @@ -458,8 +457,7 @@ async fn cleanup_strm_output_directory( let mut entries = read_dir(root_path) .await .map_err(|e| format!("Failed to read directory {root_path:?}: {e}"))?; - while let Some(entry) = entries.next().await { - let entry = entry.map_err(|e| format!("Error retrieving directory entry: {e}"))?; + while let Ok(Some(entry)) = entries.next_entry().await { if entry .file_type() .await @@ -495,7 +493,7 @@ async fn cleanup_strm_output_directory( Ok(()) } -async fn remove_empty_dirs(root_path: async_std::path::PathBuf) -> Result<(), String> { +async fn remove_empty_dirs(root_path: PathBuf) -> Result<(), String> { let mut stack = vec![root_path]; let mut dirs_to_delete = Vec::new(); let mut ignore_root = true; @@ -511,23 +509,16 @@ async fn remove_empty_dirs(root_path: async_std::path::PathBuf) -> Result<(), St let mut has_files = false; - while let Some(entry) = entries.next().await { - match entry { - Ok(entry) => { - if entry - .file_type() - .await - .map_err(|e| format!("Failed to get file type for {entry:?}: {e}"))? - .is_dir() - { - stack.push(entry.path()); - } else { - has_files = true; - } - } - Err(err) => { - error!("Error retrieving directory entry: {dir:?} {err}"); - } + while let Ok(Some(entry)) = entries.next_entry().await { + if entry + .file_type() + .await + .map_err(|e| format!("Failed to get file type for {entry:?}: {e}"))? + .is_dir() + { + stack.push(entry.path()); + } else { + has_files = true; } } @@ -539,7 +530,7 @@ async fn remove_empty_dirs(root_path: async_std::path::PathBuf) -> Result<(), St // Delete directories from bottom to top for dir in dirs_to_delete.into_iter().rev() { - if dir.exists().await { + if dir.exists() { if let Err(e) = remove_dir(&dir).await { debug!("Failed to remove empty directory {dir:?}: {e}"); } @@ -579,7 +570,7 @@ struct StrmFile { strm_info: StrmItemInfo, } -async fn prepare_strm_files( +fn prepare_strm_files( cfg: &Config, new_playlist: &[PlaylistGroup], root_path: &Path, @@ -609,7 +600,6 @@ async fn prepare_strm_files( &mut input_tmdb_indexes, underscore_whitespace, ) - .await } else { let dir_path = root_path.join(sanitize_for_filename( &strm_item_info.group, @@ -672,16 +662,14 @@ pub async fn kodi_write_strm_playlist( ))); }; - let credentials_and_server_info = get_credentials_and_server_info(cfg, output).await; + let credentials_and_server_info = get_credentials_and_server_info(cfg, output); let (underscore_whitespace, cleanup, kodi_style) = get_strm_output_options(target); let strm_index_path = strm_get_file_paths(&ensure_target_storage_path(cfg, target.name.as_str())?); let existing_strm = { let _file_lock = cfg .file_locks - .read_lock(&strm_index_path) - .await - .map_err(|err| info_err!(format!("{err}")))?; + .read_lock(&strm_index_path); read_strm_file_index(&strm_index_path) .await .unwrap_or_else(|_| HashSet::with_capacity(4096)) @@ -705,8 +693,7 @@ pub async fn kodi_write_strm_playlist( &root_path, underscore_whitespace, kodi_style, - ) - .await; + ); for strm_file in strm_files { // file paths let output_path = root_path.join(&strm_file.dir_path); @@ -737,8 +724,7 @@ pub async fn kodi_write_strm_playlist( &file_path, content_as_bytes, strm_file.strm_info.get_file_ts(), - ) - .await + ).await { Ok(()) => { processed_strm.insert(relative_file_path); @@ -772,9 +758,7 @@ async fn write_strm_index_file( ) -> Result<(), String> { let _file_lock = cfg .file_locks - .write_lock(index_file_path) - .await - .map_err(|err| format!("{err}"))?; + .write_lock(index_file_path); let file = File::create(index_file_path) .await .map_err(|err| format!("Failed to create strm index file: {index_file_path:?} {err}"))?; @@ -852,16 +836,16 @@ async fn has_strm_file_same_hash(file_path: &PathBuf, content_hash: UUIDType) -> false } -async fn get_credentials_and_server_info( +fn get_credentials_and_server_info( cfg: &Config, output: &TargetOutput, ) -> Option<(ProxyUserCredentials, ApiProxyServerInfo)> { let username = output.username.as_ref()?; - let credentials = cfg.get_user_credentials(username).await?; + let credentials = cfg.get_user_credentials(username)?; if credentials.proxy != ProxyType::Reverse { return None; } - let server_info = cfg.get_user_server_info(&credentials).await; + let server_info = cfg.get_user_server_info(&credentials); Some((credentials, server_info)) } @@ -870,7 +854,7 @@ async fn read_strm_file_index(strm_file_index_path: &Path) -> std::io::Result::new(&m3u_path, &idx_path) @@ -46,7 +45,7 @@ impl M3uPlaylistIterator { let include_type_in_url = target_options.is_some_and(|opts| opts.m3u_include_type_in_url); let mask_redirect_url = target_options.is_some_and(|opts| opts.m3u_mask_redirect_url); - let server_info = cfg.get_user_server_info(user).await; + let server_info = cfg.get_user_server_info(user); Ok(Self { reader, base_url: server_info.get_base_url(), diff --git a/src/repository/m3u_repository.rs b/src/repository/m3u_repository.rs index 7760883fb..d550a4443 100644 --- a/src/repository/m3u_repository.rs +++ b/src/repository/m3u_repository.rs @@ -3,7 +3,7 @@ use std::io::{Error, Write}; use std::path::{Path, PathBuf}; use log::error; -use crate::m3u_filter_error::{info_err, create_m3u_filter_error}; +use crate::m3u_filter_error::{create_m3u_filter_error}; use crate::m3u_filter_error::{str_to_io_error, M3uFilterError, M3uFilterErrorKind}; use crate::model::api_proxy::{ProxyUserCredentials}; use crate::model::config::{Config, ConfigTarget}; @@ -62,7 +62,7 @@ pub async fn m3u_write_playlist(target: &ConfigTarget, cfg: &Config, target_path persist_m3u_playlist_as_text(target, cfg, &m3u_playlist); { - let _file_lock = cfg.file_locks.write_lock(&m3u_path).await.map_err(|err| info_err!(format!("{err}")))?; + let _file_lock = cfg.file_locks.write_lock(&m3u_path); match IndexedDocumentWriter::new(m3u_path.clone(), idx_path) { Ok(mut writer) => { for m3u in m3u_playlist { @@ -85,7 +85,7 @@ pub async fn m3u_load_rewrite_playlist( target: &ConfigTarget, user: &ProxyUserCredentials, ) -> Result>, M3uFilterError> { - Ok(Box::new(M3uPlaylistIterator::new(cfg, target, user).await?)) + Ok(Box::new(M3uPlaylistIterator::new(cfg, target, user)?)) } @@ -96,7 +96,7 @@ pub async fn m3u_get_item_for_stream_id(stream_id: u32, cfg: &Config, target: & { let target_path = get_target_storage_path(cfg, target.name.as_str()).ok_or_else(|| str_to_io_error(&format!("Could not find path for target {}", &target.name)))?; let (m3u_path, idx_path) = m3u_get_file_paths(&target_path); - let _file_lock = cfg.file_locks.read_lock(&m3u_path).await?; + let _file_lock = cfg.file_locks.read_lock(&m3u_path); IndexedDocumentDirectAccess::read_indexed_item::(&m3u_path, &idx_path, &stream_id) } } \ No newline at end of file diff --git a/src/repository/playlist_repository.rs b/src/repository/playlist_repository.rs index fcb560708..32bcdd602 100644 --- a/src/repository/playlist_repository.rs +++ b/src/repository/playlist_repository.rs @@ -22,13 +22,7 @@ pub async fn persist_playlist(playlist: &mut [PlaylistGroup], epg: Option<&Epg>, let target_id_mapping_file = get_target_id_mapping_file(&target_path); - let _file_lock = match cfg.file_locks.write_lock(&target_id_mapping_file).await { - Ok(lock) => lock, - Err(err) => { - errors.push(info_err!(err.to_string())); - return Err(errors); - } - }; + let _file_lock = cfg.file_locks.write_lock(&target_id_mapping_file); let mut target_id_mapping = TargetIdMapping::new(&target_id_mapping_file); diff --git a/src/repository/xtream_playlist_iterator.rs b/src/repository/xtream_playlist_iterator.rs index 457a2c0da..26f27ccda 100644 --- a/src/repository/xtream_playlist_iterator.rs +++ b/src/repository/xtream_playlist_iterator.rs @@ -19,7 +19,7 @@ pub struct XtreamPlaylistIterator { } impl XtreamPlaylistIterator { - pub async fn new( + pub fn new( cluster: XtreamCluster, config: &Config, target: &ConfigTarget, @@ -31,14 +31,13 @@ impl XtreamPlaylistIterator { if !xtream_path.exists() || !idx_path.exists() { return Err(info_err!(format!("No {cluster} entries found for target {}", &target.name))); } - let file_lock = config.file_locks.read_lock(&xtream_path).await - .map_err(|err| info_err!(format!("Could not lock document {xtream_path:?}: {err}")))?; + let file_lock = config.file_locks.read_lock(&xtream_path); let reader = IndexedDocumentIterator::::new(&xtream_path, &idx_path) .map_err(|err| info_err!(format!("Could not deserialize file {xtream_path:?} - {err}")))?; let options = XtreamMappingOptions::from_target_options(target.options.as_ref(), config); - let server_info = config.get_user_server_info(user).await; + let server_info = config.get_user_server_info(user); Ok(Self { reader, options, diff --git a/src/repository/xtream_repository.rs b/src/repository/xtream_repository.rs index 37aec4673..6df2bf2c5 100644 --- a/src/repository/xtream_repository.rs +++ b/src/repository/xtream_repository.rs @@ -1,14 +1,15 @@ +use crate::m3u_filter_error::str_to_io_error; use crate::repository::storage::hex_encode; use crate::utils::file::file_utils::file_reader; -use crate::m3u_filter_error::str_to_io_error; +use log::error; +use serde_json::{json, Map, Value}; use std::collections::HashMap; use std::fs; use std::fs::File; use std::io::{BufReader, Error, ErrorKind, Read}; use std::path::{Path, PathBuf}; -use log::error; -use serde_json::{json, Map, Value}; +use crate::m3u_filter_error::{create_m3u_filter_error, create_m3u_filter_error_result, notify_err}; use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; use crate::model::api_proxy::{ProxyType, ProxyUserCredentials}; use crate::model::config::{Config, ConfigInput, ConfigTarget}; @@ -16,13 +17,12 @@ use crate::model::playlist::{PlaylistEntry, PlaylistGroup, PlaylistItem, Playlis use crate::model::xtream::{rewrite_doc_urls, XtreamMappingOptions, XtreamSeriesEpisode, INFO_RESOURCE_PREFIX, INFO_RESOURCE_PREFIX_EPISODE, SEASON_RESOURCE_PREFIX}; use crate::repository::bplustree::{BPlusTree, BPlusTreeQuery, BPlusTreeUpdate}; use crate::repository::indexed_document::{IndexedDocumentDirectAccess, IndexedDocumentGarbageCollector, IndexedDocumentWriter}; -use crate::repository::storage::{get_input_storage_path, get_target_id_mapping_file, get_target_storage_path, FILE_SUFFIX_DB, FILE_SUFFIX_INDEX}; +use crate::repository::storage::{get_input_storage_path, get_target_id_mapping_file, get_target_storage_path, FILE_SUFFIX_DB, FILE_SUFFIX_INDEX}; use crate::repository::target_id_mapping::{TargetIdMapping, VirtualIdRecord}; use crate::repository::xtream_playlist_iterator::XtreamPlaylistIterator; use crate::utils::file::file_utils::open_readonly_file; -use crate::utils::json_utils::{get_u32_from_serde_value, json_iter_array, json_write_documents_to_file}; -use crate::m3u_filter_error::{create_m3u_filter_error, create_m3u_filter_error_result, info_err, notify_err}; use crate::utils::hash_utils::generate_playlist_uuid; +use crate::utils::json_utils::{get_u32_from_serde_value, json_iter_array, json_write_documents_to_file}; pub static COL_CAT_LIVE: &str = "cat_live"; pub static COL_CAT_SERIES: &str = "cat_series"; @@ -121,7 +121,7 @@ pub fn xtream_get_record_file_path(storage_path: &Path, item_type: PlaylistItemT _ => None, } } -async fn write_playlists_to_file( +fn write_playlists_to_file( cfg: &Config, storage_path: &Path, collections: Vec<(XtreamCluster, &mut [&PlaylistItem])>, @@ -129,7 +129,7 @@ async fn write_playlists_to_file( for (cluster, playlist) in collections { let (xtream_path, idx_path) = xtream_get_file_paths(storage_path, cluster); { - let _file_lock = cfg.file_locks.write_lock(&xtream_path).await.map_err(|err| info_err!(format!("{err}")))?; + let _file_lock = cfg.file_locks.write_lock(&xtream_path); match IndexedDocumentWriter::new(xtream_path.clone(), idx_path) { Ok(mut writer) => { for item in playlist { @@ -203,7 +203,7 @@ pub fn xtream_get_file_paths_for_series(storage_path: &Path) -> (PathBuf, PathBu xtream_get_file_paths_for_name(storage_path, FILE_SERIES) } -async fn xtream_garbage_collect(config: &Config, target_name: &str) -> std::io::Result<()> { +fn xtream_garbage_collect(config: &Config, target_name: &str) -> std::io::Result<()> { // Garbage collect series let storage_path = try_option_ok!(xtream_get_storage_path(config, target_name)); let (info_path, idx_path) = try_option_ok!(xtream_get_info_file_paths( @@ -211,7 +211,7 @@ async fn xtream_garbage_collect(config: &Config, target_name: &str) -> std::io:: XtreamCluster::Series )); { - let _file_lock = config.file_locks.write_lock(&info_path).await?; + let _file_lock = config.file_locks.write_lock(&info_path); IndexedDocumentGarbageCollector::::new(info_path, idx_path)?.garbage_collect()?; } Ok(()) @@ -263,25 +263,25 @@ pub async fn xtream_write_playlist( }; // let col = match header.item_type { - // PlaylistItemType::LiveUnknown | PlaylistItemType::LiveHls => { - // header.category_id = *cat_id; - // Some(&mut live_col) - // } - // _ => { - // if header.get_provider_id().is_some() { - // header.category_id = *cat_id; - // Some(match header.xtream_cluster { - // XtreamCluster::Live => &mut live_col, - // XtreamCluster::Series => &mut series_col, - // XtreamCluster::Video => &mut vod_col, - // }) - // // } else { - // let title = header.title.as_str(); - // errors.push(format!("Channel does not have an id: {title}")); - // errors.push(format!("Channel does not have an id: {title}")); - // None - // } - // } + // PlaylistItemType::LiveUnknown | PlaylistItemType::LiveHls => { + // header.category_id = *cat_id; + // Some(&mut live_col) + // } + // _ => { + // if header.get_provider_id().is_some() { + // header.category_id = *cat_id; + // Some(match header.xtream_cluster { + // XtreamCluster::Live => &mut live_col, + // XtreamCluster::Series => &mut series_col, + // XtreamCluster::Video => &mut vod_col, + // }) + // // } else { + // let title = header.title.as_str(); + // errors.push(format!("Channel does not have an id: {title}")); + // errors.push(format!("Channel does not have an id: {title}")); + // None + // } + // } // }; drop(header); col.push(pli); @@ -310,9 +310,9 @@ pub async fn xtream_write_playlist( (XtreamCluster::Video, &mut vod_col), (XtreamCluster::Series, &mut series_col), ], - ).await { + ) { Ok(()) => { - if let Err(err) = xtream_garbage_collect(cfg, &target.name).await { + if let Err(err) = xtream_garbage_collect(cfg, &target.name) { if err.kind() != ErrorKind::NotFound { errors.push(format!("Garbage collection failed:{err}")); } @@ -348,7 +348,7 @@ pub fn xtream_get_collection_path( Err(str_to_io_error(&format!("Cant find collection: {target_name}/{collection_name}"))) } -async fn xtream_read_item_for_stream_id( +fn xtream_read_item_for_stream_id( cfg: &Config, stream_id: u32, storage_path: &Path, @@ -356,19 +356,19 @@ async fn xtream_read_item_for_stream_id( ) -> Result { let (xtream_path, idx_path) = xtream_get_file_paths(storage_path, cluster); { - let _file_lock = cfg.file_locks.read_lock(&xtream_path).await?; + let _file_lock = cfg.file_locks.read_lock(&xtream_path); IndexedDocumentDirectAccess::read_indexed_item::(&xtream_path, &idx_path, &stream_id) } } -async fn xtream_read_series_item_for_stream_id( +fn xtream_read_series_item_for_stream_id( cfg: &Config, stream_id: u32, storage_path: &Path, ) -> Result { let (xtream_path, idx_path) = xtream_get_file_paths_for_series(storage_path); { - let _file_lock = cfg.file_locks.read_lock(&xtream_path).await?; + let _file_lock = cfg.file_locks.read_lock(&xtream_path); IndexedDocumentDirectAccess::read_indexed_item::(&xtream_path, &idx_path, &stream_id) } } @@ -381,7 +381,7 @@ macro_rules! try_cluster { }; } -pub async fn xtream_get_item_for_stream_id( +pub fn xtream_get_item_for_stream_id( virtual_id: u32, config: &Config, target: &ConfigTarget, @@ -391,31 +391,31 @@ pub async fn xtream_get_item_for_stream_id( let storage_path = xtream_get_storage_path(config, target.name.as_str()).ok_or_else(|| str_to_io_error(&format!("Could not find path for target {} xtream output", &target.name)))?; { let target_id_mapping_file = get_target_id_mapping_file(&target_path); - let _file_lock = config.file_locks.read_lock(&target_id_mapping_file).await.map_err(|err| str_to_io_error(&format!("Could not get lock for id mapping for target {} err:{err}", target.name)))?; + let _file_lock = config.file_locks.read_lock(&target_id_mapping_file); let mut target_id_mapping = BPlusTreeQuery::::try_new(&target_id_mapping_file).map_err(|err| str_to_io_error(&format!("Could not load id mapping for target {} err:{err}", target.name)))?; let mapping = target_id_mapping.query(&virtual_id).ok_or_else(|| str_to_io_error(&format!("Could not find mapping for target {} and id {}", target.name, virtual_id)))?; let result = match mapping.item_type { PlaylistItemType::SeriesInfo => { - xtream_read_series_item_for_stream_id(config, virtual_id, &storage_path).await + xtream_read_series_item_for_stream_id(config, virtual_id, &storage_path) } PlaylistItemType::Series => { - if let Ok(mut item) = xtream_read_series_item_for_stream_id(config, mapping.parent_virtual_id, &storage_path).await { + if let Ok(mut item) = xtream_read_series_item_for_stream_id(config, mapping.parent_virtual_id, &storage_path) { item.provider_id = mapping.provider_id; Ok(item) } else { - xtream_read_item_for_stream_id(config, virtual_id, &storage_path, XtreamCluster::Series).await + xtream_read_item_for_stream_id(config, virtual_id, &storage_path, XtreamCluster::Series) } } PlaylistItemType::Catchup => { let cluster = try_cluster!(xtream_cluster, mapping.item_type, virtual_id)?; - let mut item = xtream_read_item_for_stream_id(config, mapping.parent_virtual_id, &storage_path, cluster).await?; + let mut item = xtream_read_item_for_stream_id(config, mapping.parent_virtual_id, &storage_path, cluster)?; item.provider_id = mapping.provider_id; Ok(item) } _ => { let cluster = try_cluster!(xtream_cluster, mapping.item_type, virtual_id)?; - xtream_read_item_for_stream_id(config, virtual_id, &storage_path, cluster).await + xtream_read_item_for_stream_id(config, virtual_id, &storage_path, cluster) } }; @@ -423,17 +423,17 @@ pub async fn xtream_get_item_for_stream_id( } } -pub async fn xtream_load_rewrite_playlist( +pub fn xtream_load_rewrite_playlist( cluster: XtreamCluster, config: &Config, target: &ConfigTarget, category_id: u32, user: &ProxyUserCredentials, ) -> Result>, M3uFilterError> { - Ok(Box::new(XtreamPlaylistIterator::new(cluster, config, target, category_id, user).await?)) + Ok(Box::new(XtreamPlaylistIterator::new(cluster, config, target, category_id, user)?)) } -pub async fn xtream_write_series_info( +pub fn xtream_write_series_info( config: &Config, target_name: &str, series_info_id: u32, @@ -447,14 +447,14 @@ pub async fn xtream_write_series_info( )); { - let _file_lock = config.file_locks.write_lock(&info_path).await?; + let _file_lock = config.file_locks.write_lock(&info_path); let mut writer = IndexedDocumentWriter::new_append(info_path, idx_path)?; writer.write_doc(series_info_id, content).map_err(|_| str_to_io_error(&format!("failed to write xtream series info for target {target_name}")))?; writer.store()?; } { let target_id_mapping_file = get_target_id_mapping_file(&target_path); - let _file_lock = config.file_locks.write_lock(&target_id_mapping_file).await?; + let _file_lock = config.file_locks.write_lock(&target_id_mapping_file); if let Ok(mut target_id_mapping) = BPlusTreeUpdate::::try_new(&target_id_mapping_file) { if let Some(record) = target_id_mapping.query(&series_info_id) { let new_record = record.copy_update_timestamp(); @@ -466,7 +466,7 @@ pub async fn xtream_write_series_info( Ok(()) } -pub async fn xtream_write_vod_info( +pub fn xtream_write_vod_info( config: &Config, target_name: &str, virtual_id: u32, @@ -475,7 +475,7 @@ pub async fn xtream_write_vod_info( let storage_path = try_option_ok!(xtream_get_storage_path(config, target_name)); let (info_path, idx_path) = try_option_ok!(xtream_get_info_file_paths(&storage_path, XtreamCluster::Video)); { - let _file_lock = config.file_locks.write_lock(&info_path).await?; + let _file_lock = config.file_locks.write_lock(&info_path); let mut writer = IndexedDocumentWriter::new_append(info_path, idx_path)?; writer.write_doc(virtual_id, content).map_err(|_| str_to_io_error(&format!("failed to write xtream vod info for target {target_name}")))?; writer.store()?; @@ -483,22 +483,19 @@ pub async fn xtream_write_vod_info( Ok(()) } -async fn xtream_get_series_info_mapping( +fn xtream_get_series_info_mapping( config: &Config, target_name: &str, series_id: u32, ) -> Option { - xtream_get_info_mapping(config, target_name, series_id).await.filter(|id_record| !id_record.is_expired()) + xtream_get_info_mapping(config, target_name, series_id).filter(|id_record| !id_record.is_expired()) } -async fn xtream_get_info_mapping(config: &Config, target_name: &str, info_id: u32) -> Option { +fn xtream_get_info_mapping(config: &Config, target_name: &str, info_id: u32) -> Option { let target_path = get_target_storage_path(config, target_name)?; let target_id_mapping_file = get_target_id_mapping_file(&target_path); - let _file_lock = config.file_locks.read_lock(&target_id_mapping_file).await.map_err(|err| { - error!("Could not lock id mapping for target {target_name}: {}", err); - str_to_io_error(&format!("ID mapping load error for target {target_name}")) - }).ok()?; + let _file_lock = config.file_locks.read_lock(&target_id_mapping_file); BPlusTreeQuery::::try_new(&target_id_mapping_file).map_err(|err| { error!("Could not load id mapping for target {target_name}: {}", err); str_to_io_error(&format!("ID mapping load error for target {target_name}")) @@ -506,12 +503,12 @@ async fn xtream_get_info_mapping(config: &Config, target_name: &str, info_id: u3 } // Reads the series info entry if exists -pub async fn xtream_load_series_info( +pub fn xtream_load_series_info( config: &Config, target_name: &str, series_id: u32, ) -> Option { - xtream_get_series_info_mapping(config, target_name, series_id).await?; + xtream_get_series_info_mapping(config, target_name, series_id)?; let storage_path = xtream_get_storage_path(config, target_name)?; @@ -519,10 +516,7 @@ pub async fn xtream_load_series_info( if info_path.exists() && idx_path.exists() { { - let _file_lock = config.file_locks.read_lock(&info_path).await.map_err(|err| { - error!("Could not lock document {:?}: {}", info_path, err); - str_to_io_error(&format!("Document Reader error for target {target_name}")) - }).ok()?; + let _file_lock = config.file_locks.read_lock(&info_path); return match IndexedDocumentDirectAccess::read_indexed_item::(&info_path, &idx_path, &series_id) { Ok(content) => Some(content), Err(err) => { @@ -534,25 +528,24 @@ pub async fn xtream_load_series_info( } None } - -async fn xtream_get_vod_info_mapping( + fn xtream_get_vod_info_mapping( config: &Config, target_name: &str, vod_id: u32, ) -> Option { - xtream_get_info_mapping(config, target_name, vod_id).await + xtream_get_info_mapping(config, target_name, vod_id) //.filter(|id_record| !id_record.is_expired()) } // Reads the vod info entry if exists -pub async fn xtream_load_vod_info( +pub fn xtream_load_vod_info( config: &Config, target_name: &str, vod_id: u32, ) -> Option { // Check if the entry exists; if not, we don't need to look further. - xtream_get_vod_info_mapping(config, target_name, vod_id).await.as_ref()?; + xtream_get_vod_info_mapping(config, target_name, vod_id).as_ref()?; // Entry exists, read db entry let target_storage_path = xtream_get_storage_path(config, target_name)?; @@ -560,10 +553,7 @@ pub async fn xtream_load_vod_info( if info_path.exists() && idx_path.exists() { { - let _file_lock = config.file_locks.read_lock(&info_path).await.map_err(|err| { - error!("Could not lock document {:?}: {}", info_path, err); - str_to_io_error(&format!("Document Reader error for target {target_name}")) - }).ok()?; + let _file_lock = config.file_locks.read_lock(&info_path); return match IndexedDocumentDirectAccess::read_indexed_item::(&info_path, &idx_path, &vod_id) { Ok(content) => Some(content), Err(_err) => { @@ -577,7 +567,7 @@ pub async fn xtream_load_vod_info( None } -async fn rewrite_xtream_vod_info

( +fn rewrite_xtream_vod_info

( config: &Config, target: &ConfigTarget, pli: &P, @@ -591,7 +581,7 @@ async fn rewrite_xtream_vod_info

( if let Some(Value::Object(info_data)) = doc.get_mut(TAG_INFO_DATA) { match user.proxy { ProxyType::Reverse => { - let server_info = config.get_user_server_info(user).await; + let server_info = config.get_user_server_info(user); let url = server_info.get_base_url(); let resource_url = Some(format!("{url}/resource/movie/{}/{}/{}", user.username, user.password, pli.get_virtual_id())); rewrite_doc_urls(resource_url.as_ref(), info_data, INFO_REWRITE_FIELDS, INFO_RESOURCE_PREFIX); @@ -624,7 +614,7 @@ async fn rewrite_xtream_vod_info

( Ok(result) } -pub async fn rewrite_xtream_vod_info_content

( +pub fn rewrite_xtream_vod_info_content

( config: &Config, target: &ConfigTarget, pli: &P, @@ -634,10 +624,10 @@ pub async fn rewrite_xtream_vod_info_content

( P: PlaylistEntry, { let mut doc = serde_json::from_str::>(content).map_err(|_| str_to_io_error("Failed to parse JSON content"))?; - rewrite_xtream_vod_info(config, target, pli, user, &mut doc).await + rewrite_xtream_vod_info(config, target, pli, user, &mut doc) } -pub async fn write_and_get_xtream_vod_info

( +pub fn write_and_get_xtream_vod_info

( config: &Config, target: &ConfigTarget, pli: &P, @@ -647,11 +637,11 @@ pub async fn write_and_get_xtream_vod_info

( P: PlaylistEntry, { let mut doc = serde_json::from_str::>(content).map_err(|_| str_to_io_error("Failed to parse JSON content"))?; - xtream_write_vod_info(config, target.name.as_str(), pli.get_virtual_id(), content).await.ok(); - rewrite_xtream_vod_info(config, target, pli, user, &mut doc).await + xtream_write_vod_info(config, target.name.as_str(), pli.get_virtual_id(), content).ok(); + rewrite_xtream_vod_info(config, target, pli, user, &mut doc) } -async fn rewrite_xtream_series_info

( +fn rewrite_xtream_series_info

( config: &Config, target: &ConfigTarget, pli: &P, @@ -665,7 +655,7 @@ async fn rewrite_xtream_series_info

( let resource_url = if config.is_reverse_proxy_resource_rewrite_enabled() { match user.proxy { ProxyType::Reverse => { - let server_info = config.get_user_server_info(user).await; + let server_info = config.get_user_server_info(user); let url = server_info.get_base_url(); Some(format!("{url}/resource/series/{}/{}/{}", user.username, user.password, pli.get_virtual_id())) } @@ -696,7 +686,7 @@ async fn rewrite_xtream_series_info

( let virtual_id = pli.get_virtual_id(); { let target_id_mapping_file = get_target_id_mapping_file(&target_path); - let _file_lock = config.file_locks.write_lock(&target_id_mapping_file).await.map_err(|err| str_to_io_error(&format!("Could not load id mapping for target {} err:{err}", target.name)))?; + let _file_lock = config.file_locks.write_lock(&target_id_mapping_file); let mut target_id_mapping = TargetIdMapping::new(&target_id_mapping_file); let options = XtreamMappingOptions::from_target_options(target.options.as_ref(), config); @@ -737,7 +727,7 @@ async fn rewrite_xtream_series_info

( Ok(result) } -pub async fn rewrite_xtream_series_info_content

( +pub fn rewrite_xtream_series_info_content

( config: &Config, target: &ConfigTarget, pli_series_info: &P, @@ -747,10 +737,10 @@ pub async fn rewrite_xtream_series_info_content

( P: PlaylistEntry, { let mut doc = serde_json::from_str::>(content).map_err(|_| str_to_io_error("Failed to parse JSON content"))?; - rewrite_xtream_series_info(config, target, pli_series_info, user, &mut doc).await + rewrite_xtream_series_info(config, target, pli_series_info, user, &mut doc) } -pub async fn write_and_get_xtream_series_info

( +pub fn write_and_get_xtream_series_info

( config: &Config, target: &ConfigTarget, pli_series_info: &P, @@ -761,11 +751,11 @@ pub async fn write_and_get_xtream_series_info

( { let mut doc = serde_json::from_str::>(content).map_err(|_| str_to_io_error("Failed to parse JSON content"))?; let virtual_id = pli_series_info.get_virtual_id(); - xtream_write_series_info(config, target.name.as_str(), virtual_id, content).await.ok(); - rewrite_xtream_series_info(config, target, pli_series_info, user, &mut doc).await + xtream_write_series_info(config, target.name.as_str(), virtual_id, content).ok(); + rewrite_xtream_series_info(config, target, pli_series_info, user, &mut doc) } -pub async fn xtream_get_input_info( +pub fn xtream_get_input_info( cfg: &Config, input: &ConfigInput, provider_id: u32, @@ -773,10 +763,9 @@ pub async fn xtream_get_input_info( ) -> Option { if let Ok(Some((info_path, idx_path))) = get_input_storage_path(input, &cfg.working_dir).map(|storage_path| xtream_get_info_file_paths(&storage_path, cluster)) { - if let Ok(_file_lock) = cfg.file_locks.read_lock(&info_path).await { - if let Ok(content) = IndexedDocumentDirectAccess::read_indexed_item::(&info_path, &idx_path, &provider_id) { - return Some(content); - } + let _file_lock = cfg.file_locks.read_lock(&info_path); + if let Ok(content) = IndexedDocumentDirectAccess::read_indexed_item::(&info_path, &idx_path, &provider_id) { + return Some(content); } } None @@ -790,37 +779,35 @@ pub async fn xtream_update_input_info_file( ) -> Result<(), M3uFilterError> { match get_input_storage_path(input, &cfg.working_dir).map(|storage_path| xtream_get_info_file_paths(&storage_path, cluster)) { Ok(Some((info_path, idx_path))) => { - match cfg.file_locks.write_lock(&info_path).await { - Ok(_file_lock) => { - let mut reader = file_reader(open_readonly_file(wal_path).map_err(|err| notify_err!(format!("Could not read {cluster} info {err}")))?); - match IndexedDocumentWriter::::new_append(info_path, idx_path) { - Ok(mut writer) => { - let mut provider_id_bytes = [0u8; 4]; - let mut length_bytes = [0u8; 4]; - loop { - if reader.read_exact(&mut provider_id_bytes).is_err() { - break; // End of file - } - let provider_id = u32::from_le_bytes(provider_id_bytes); - reader.read_exact(&mut length_bytes).map_err(|err| notify_err!(format!("Could not read temporary {cluster} info {err}")))?; - let length = u32::from_le_bytes(length_bytes) as usize; - let mut buffer = vec![0u8; length]; - reader.read_exact(&mut buffer).map_err(|err| notify_err!(format!("Could not read temporary {cluster} info {err}")))?; - if let Ok(content) = String::from_utf8(buffer) { - let _ = writer.write_doc(provider_id, &content); - } + { + let _file_lock = cfg.file_locks.write_lock(&info_path); + let mut reader = file_reader(open_readonly_file(wal_path).map_err(|err| notify_err!(format!("Could not read {cluster} info {err}")))?); + match IndexedDocumentWriter::::new_append(info_path, idx_path) { + Ok(mut writer) => { + let mut provider_id_bytes = [0u8; 4]; + let mut length_bytes = [0u8; 4]; + loop { + if reader.read_exact(&mut provider_id_bytes).is_err() { + break; // End of file } - writer.store().map_err(|err| notify_err!(format!("Could not store {cluster} info {err}")))?; - drop(reader); - if let Err(err) = fs::remove_file(wal_path) { - error!("Failed to delete WAL file for {cluster} {err}"); + let provider_id = u32::from_le_bytes(provider_id_bytes); + reader.read_exact(&mut length_bytes).map_err(|err| notify_err!(format!("Could not read temporary {cluster} info {err}")))?; + let length = u32::from_le_bytes(length_bytes) as usize; + let mut buffer = vec![0u8; length]; + reader.read_exact(&mut buffer).map_err(|err| notify_err!(format!("Could not read temporary {cluster} info {err}")))?; + if let Ok(content) = String::from_utf8(buffer) { + let _ = writer.write_doc(provider_id, &content); } - Ok(()) } - Err(err) => Err(notify_err!(format!("Could not create create indexed document writer for {cluster} info {err}"))), + writer.store().map_err(|err| notify_err!(format!("Could not store {cluster} info {err}")))?; + drop(reader); + if let Err(err) = fs::remove_file(wal_path) { + error!("Failed to delete WAL file for {cluster} {err}"); + } + Ok(()) } + Err(err) => Err(notify_err!(format!("Could not create create indexed document writer for {cluster} info {err}"))), } - Err(err) => Err(info_err!(format!("{err}"))), } } Ok(None) => Err(notify_err!(format!("Could not create storage path for input {}", &input.name))), @@ -837,36 +824,34 @@ pub async fn xtream_update_input_vod_record_from_wal_file( .map_err(|err| notify_err!(format!("Error accessing storage path: {err}"))) .and_then(|opt| opt.ok_or_else(|| notify_err!(format!("Error accessing storage path for input: {}", &input.name))))?; - match cfg.file_locks.write_lock(&record_path).await { - Ok(_file_lock) => { - let mut reader = file_reader(open_readonly_file(wal_path).map_err(|err| notify_err!(format!("Could not read vod wal info {err}")))?); - let mut provider_id_bytes = [0u8; 4]; - let mut tmdb_id_bytes = [0u8; 4]; - let mut ts_bytes = [0u8; 8]; - let mut tree_record_index: BPlusTree = BPlusTree::load(&record_path).unwrap_or_else(|_| BPlusTree::new()); - loop { - if reader.read_exact(&mut provider_id_bytes).is_err() { - break; // End of file - } - let provider_id = u32::from_le_bytes(provider_id_bytes); - if reader.read_exact(&mut tmdb_id_bytes).is_err() { - break; // End of file - } - let tmdb_id = u32::from_le_bytes(tmdb_id_bytes); - if reader.read_exact(&mut ts_bytes).is_err() { - break; // End of file - } - let ts = u64::from_le_bytes(ts_bytes); - tree_record_index.insert(provider_id, InputVodInfoRecord { tmdb_id, ts }); + { + let _file_lock = cfg.file_locks.write_lock(&record_path); + let mut reader = file_reader(open_readonly_file(wal_path).map_err(|err| notify_err!(format!("Could not read vod wal info {err}")))?); + let mut provider_id_bytes = [0u8; 4]; + let mut tmdb_id_bytes = [0u8; 4]; + let mut ts_bytes = [0u8; 8]; + let mut tree_record_index: BPlusTree = BPlusTree::load(&record_path).unwrap_or_else(|_| BPlusTree::new()); + loop { + if reader.read_exact(&mut provider_id_bytes).is_err() { + break; // End of file } - tree_record_index.store(&record_path).map_err(|err| notify_err!(format!("Could not store vod record info {err}")))?; - drop(reader); - if let Err(err) = fs::remove_file(wal_path) { - error!("Failed to delete record WAL file for vod {err}"); + let provider_id = u32::from_le_bytes(provider_id_bytes); + if reader.read_exact(&mut tmdb_id_bytes).is_err() { + break; // End of file } - Ok(()) + let tmdb_id = u32::from_le_bytes(tmdb_id_bytes); + if reader.read_exact(&mut ts_bytes).is_err() { + break; // End of file + } + let ts = u64::from_le_bytes(ts_bytes); + tree_record_index.insert(provider_id, InputVodInfoRecord { tmdb_id, ts }); } - Err(err) => Err(info_err!(format!("{err}"))), + tree_record_index.store(&record_path).map_err(|err| notify_err!(format!("Could not store vod record info {err}")))?; + drop(reader); + if let Err(err) = fs::remove_file(wal_path) { + error!("Failed to delete record WAL file for vod {err}"); + } + Ok(()) } } @@ -878,32 +863,29 @@ pub async fn xtream_update_input_series_record_from_wal_file( let record_path = get_input_storage_path(input, &cfg.working_dir).map(|storage_path| xtream_get_record_file_path(&storage_path, PlaylistItemType::SeriesInfo)) .map_err(|err| notify_err!(format!("Error accessing storage path: {err}"))) .and_then(|opt| opt.ok_or_else(|| notify_err!(format!("Error accessing storage path for input: {}", &input.name))))?; - match cfg.file_locks.write_lock(&record_path).await { - Ok(_file_lock) => { - let mut reader = file_reader(open_readonly_file(wal_path).map_err(|err| notify_err!(format!("Could not read series wal info {err}")))?); - let mut provider_id_bytes = [0u8; 4]; - let mut ts_bytes = [0u8; 8]; - let mut tree_record_index: BPlusTree = BPlusTree::load(&record_path).unwrap_or_else(|_| BPlusTree::new()); - loop { - if reader.read_exact(&mut provider_id_bytes).is_err() { - break; // End of file - } - let provider_id = u32::from_le_bytes(provider_id_bytes); - if reader.read_exact(&mut ts_bytes).is_err() { - break; // End of file - } - let ts = u64::from_le_bytes(ts_bytes); - tree_record_index.insert(provider_id, ts); + { + let _file_lock = cfg.file_locks.write_lock(&record_path); + let mut reader = file_reader(open_readonly_file(wal_path).map_err(|err| notify_err!(format!("Could not read series wal info {err}")))?); + let mut provider_id_bytes = [0u8; 4]; + let mut ts_bytes = [0u8; 8]; + let mut tree_record_index: BPlusTree = BPlusTree::load(&record_path).unwrap_or_else(|_| BPlusTree::new()); + loop { + if reader.read_exact(&mut provider_id_bytes).is_err() { + break; // End of file } - tree_record_index.store(&record_path).map_err(|err| notify_err!(format!("Could not store series record info {err}")))?; - drop(reader); - if let Err(err) = fs::remove_file(wal_path) { - error!("Failed to delete record WAL file for series {err}"); + let provider_id = u32::from_le_bytes(provider_id_bytes); + if reader.read_exact(&mut ts_bytes).is_err() { + break; // End of file } - Ok(()) + let ts = u64::from_le_bytes(ts_bytes); + tree_record_index.insert(provider_id, ts); } - - Err(err) => Err(info_err!(format!("{err}"))), + tree_record_index.store(&record_path).map_err(|err| notify_err!(format!("Could not store series record info {err}")))?; + drop(reader); + if let Err(err) = fs::remove_file(wal_path) { + error!("Failed to delete record WAL file for series {err}"); + } + Ok(()) } } @@ -915,49 +897,46 @@ pub async fn xtream_update_input_series_episodes_record_from_wal_file( let record_path = get_input_storage_path(input, &cfg.working_dir).map(|storage_path| xtream_get_record_file_path(&storage_path, PlaylistItemType::Series)) .map_err(|err| notify_err!(format!("Error accessing storage path: {err}"))) .and_then(|opt| opt.ok_or_else(|| notify_err!(format!("Error accessing storage path for input: {}", &input.name))))?; - match cfg.file_locks.write_lock(&record_path).await { - Ok(_file_lock) => { - let mut reader = file_reader(open_readonly_file(wal_path).map_err(|err| notify_err!(format!("Could not read series episode wal info {err}")))?); - let mut provider_id_bytes = [0u8; 4]; - let mut len_bytes = [0u8; 4]; - let mut tree_record_index: BPlusTree = BPlusTree::load(&record_path).unwrap_or_else(|_| BPlusTree::new()); - let mut buffer = vec![0u8; 4096]; - loop { - if reader.read_exact(&mut provider_id_bytes).is_err() { - break; // End of file + { + let _file_lock = cfg.file_locks.write_lock(&record_path); + let mut reader = file_reader(open_readonly_file(wal_path).map_err(|err| notify_err!(format!("Could not read series episode wal info {err}")))?); + let mut provider_id_bytes = [0u8; 4]; + let mut len_bytes = [0u8; 4]; + let mut tree_record_index: BPlusTree = BPlusTree::load(&record_path).unwrap_or_else(|_| BPlusTree::new()); + let mut buffer = vec![0u8; 4096]; + loop { + if reader.read_exact(&mut provider_id_bytes).is_err() { + break; // End of file + } + let provider_id = u32::from_le_bytes(provider_id_bytes); + if reader.read_exact(&mut len_bytes).is_err() { + break; // End of file + } + let len = usize::try_from(u32::from_le_bytes(len_bytes)).unwrap_or(0); + if len == 0 { + break; + } + if len > buffer.len() { + buffer = vec![0u8; len]; + } + if reader.read_exact(&mut buffer[0..len]).is_err() { + break; + } + match bincode::deserialize(&buffer[0..len]) { + Ok(episode) => { + tree_record_index.insert(provider_id, episode); } - let provider_id = u32::from_le_bytes(provider_id_bytes); - if reader.read_exact(&mut len_bytes).is_err() { - break; // End of file - } - let len = usize::try_from(u32::from_le_bytes(len_bytes)).unwrap_or(0); - if len == 0 { - break; - } - if len > buffer.len() { - buffer = vec![0u8; len]; - } - if reader.read_exact(&mut buffer[0..len]).is_err() { - break; - } - match bincode::deserialize(&buffer[0..len]) { - Ok(episode) => { - tree_record_index.insert(provider_id, episode); - } - Err(err) => { - error!("Failed to delete deserialize record WAL file for series episode {err}"); - } + Err(err) => { + error!("Failed to delete deserialize record WAL file for series episode {err}"); } } - tree_record_index.store(&record_path).map_err(|err| notify_err!(format!("Could not store series episode record info {err}")))?; - drop(reader); - if let Err(err) = fs::remove_file(wal_path) { - error!("Failed to delete record WAL file for series episode {err}"); - } - Ok(()) } - - Err(err) => Err(info_err!(format!("{err}"))), + tree_record_index.store(&record_path).map_err(|err| notify_err!(format!("Could not store series episode record info {err}")))?; + drop(reader); + if let Err(err) = fs::remove_file(wal_path) { + error!("Failed to delete record WAL file for series episode {err}"); + } + Ok(()) } } diff --git a/src/tools/atomic_once_flag.rs b/src/tools/atomic_once_flag.rs index ba5133f29..982e4c49d 100644 --- a/src/tools/atomic_once_flag.rs +++ b/src/tools/atomic_once_flag.rs @@ -1,5 +1,4 @@ use std::sync::atomic::{AtomicBool, Ordering}; -use std::sync::Arc; /// A flag that is initially active (`true`) and can only be disabled once. /// Once the flag is disabled by calling `disable()`, it remains inactive (`false`) forever. @@ -14,9 +13,9 @@ use std::sync::Arc; /// assert!(flag.is_active()); /// flag.disable(); /// assert!(!flag.is_active()); -#[derive(Clone, Debug)] +#[derive(Debug)] pub struct AtomicOnceFlag { - enabled: Arc, + enabled: AtomicBool, ordering: Ordering, } @@ -30,7 +29,7 @@ impl AtomicOnceFlag { /// Creates a new `AtomicOnceFlag` with the specified memory ordering. pub fn with_ordering(ordering: Ordering) -> Self { Self { - enabled: Arc::new(AtomicBool::new(true)), + enabled: AtomicBool::new(true), ordering, } } diff --git a/src/tools/lru_cache.rs b/src/tools/lru_cache.rs index 36b3ba170..662f54587 100644 --- a/src/tools/lru_cache.rs +++ b/src/tools/lru_cache.rs @@ -1,7 +1,7 @@ use crate::repository::storage::hash_string_as_hex; use crate::utils::file::file_utils::traverse_dir; use crate::utils::size_utils::human_readable_byte_size; -use async_std::sync::RwLock; +use parking_lot::RwLock; use log::{debug, error, info, trace}; use std::collections::{HashMap, VecDeque}; use std::fs; @@ -49,8 +49,8 @@ impl LRUResourceCache { /// - Scans the cache directory and populates the internal data structures with existing files and their sizes. /// - Updates the `current_size` and `usage_order` fields based on the scanned files. /// The use/access order is not restored!!! - pub async fn scan(&mut self) -> std::io::Result<()> { - let _write_lock = self.lock.write().await; + pub fn scan(&mut self) -> std::io::Result<()> { + let _write_lock = self.lock.write(); let mut visit = |entry: &std::fs::DirEntry, metadata: &std::fs::Metadata| { let path = entry.path(); if let Some(file_name) = path.file_name() { @@ -79,19 +79,17 @@ impl LRUResourceCache { /// - `file_size`: The size of the file in bytes. /// - Returns: /// - The `PathBuf` where the file is stored. - pub async fn add_content(&mut self, url: &str, file_size: usize) -> std::io::Result { + pub fn add_content(&mut self, url: &str, file_size: usize) -> std::io::Result { let key = hash_string_as_hex(url); - let path = { - self.insert_to_cache(key, file_size).await - }; + let path = self.insert_to_cache(key, file_size); if self.current_size > self.capacity { - self.evict_if_needed().await; + self.evict_if_needed(); } Ok(path) } - async fn insert_to_cache(&mut self, key: String, file_size: usize) -> PathBuf { - let _write_lock = self.lock.write().await; + fn insert_to_cache(&mut self, key: String, file_size: usize) -> PathBuf { + let _write_lock = self.lock.write(); let mut path = self.cache_dir.clone(); path.push(&key); debug!("Added file to cache: {}", &path.to_string_lossy()); @@ -114,10 +112,10 @@ impl LRUResourceCache { /// - `url`: The unique identifier for the file. /// - Returns: /// - The `PathBuf` of the file if it exists; `None` otherwise. - pub async fn get_content(&mut self, url: &str) -> Option { + pub fn get_content(&mut self, url: &str) -> Option { let key = hash_string_as_hex(url); { - let _read_lock = self.lock.read().await; + let _read_lock = self.lock.read(); if let Some((path, size)) = self.cache.get(&key) { if path.exists() { // Move to the end of the queue @@ -127,7 +125,7 @@ impl LRUResourceCache { } { // this should not happen, someone deleted the file manually and the cache is not in sync - let _write_lock = self.lock.write().await; + let _write_lock = self.lock.write(); self.current_size -= size; self.cache.remove(&key); self.usage_order.retain(|k| k != &key); @@ -137,8 +135,8 @@ impl LRUResourceCache { None } - async fn evict_if_needed(&mut self) { - let _write_lock = self.lock.write().await; + fn evict_if_needed(&mut self) { + let _write_lock = self.lock.write(); // if the cache size is to small and one element exceeds the size than the cache won't work, we ignore this while self.current_size > self.capacity { if let Some(oldest_file) = self.usage_order.pop_front() { diff --git a/src/utils/config_reader.rs b/src/utils/file/config_reader.rs similarity index 99% rename from src/utils/config_reader.rs rename to src/utils/file/config_reader.rs index 2f8a55a82..1e9dea170 100644 --- a/src/utils/config_reader.rs +++ b/src/utils/file/config_reader.rs @@ -152,7 +152,7 @@ pub fn resolve_env_var(value: &str) -> String { #[cfg(test)] mod tests { - use crate::utils::config_reader::resolve_env_var; + use crate::utils::file::config_reader::resolve_env_var; #[test] fn test_resolve() { diff --git a/src/utils/file/file_lock_manager.rs b/src/utils/file/file_lock_manager.rs index 7968836c4..280b362c2 100644 --- a/src/utils/file/file_lock_manager.rs +++ b/src/utils/file/file_lock_manager.rs @@ -1,8 +1,8 @@ use std::collections::HashMap; -use std::sync::{Arc}; +use std::sync::Arc; use std::{fmt, io}; use std::path::{Path, PathBuf}; -use async_std::sync::{Mutex, RwLock, RwLockReadGuard, RwLockWriteGuard}; +use parking_lot::{Mutex, RwLock, RwLockReadGuard, RwLockWriteGuard}; use crate::m3u_filter_error::str_to_io_error; #[derive(Clone)] @@ -18,24 +18,24 @@ impl FileLockManager { } // Acquires a read lock for the specified file and returns a FileReadGuard. - pub async fn read_lock(&self, path: &Path) -> io::Result { - let file_lock = self.get_or_create_lock(path).await.map_err(|_| str_to_io_error("Failed to acquire write lock"))?; - let guard = file_lock.read().await; + pub fn read_lock(&self, path: &Path) -> FileReadGuard { + let file_lock = self.get_or_create_lock(path); + let guard = file_lock.read(); // Clone the Arc to avoid moving `file_lock` out, as it is still borrowed by `guard` - Ok(FileReadGuard::new(Arc::clone(&file_lock), guard)) + FileReadGuard::new(Arc::clone(&file_lock), guard) } // Acquires a write lock for the specified file and returns a FileWriteGuard. - pub async fn write_lock(&self, path: &Path) -> io::Result { - let file_lock = self.get_or_create_lock(path).await.map_err(|_| str_to_io_error("Failed to acquire write lock"))?; - let guard = file_lock.write().await; + pub fn write_lock(&self, path: &Path) -> FileWriteGuard { + let file_lock = self.get_or_create_lock(path); + let guard = file_lock.write(); // Clone the Arc to avoid moving `file_lock` out, as it is still borrowed by `guard` - Ok(FileWriteGuard::new(Arc::clone(&file_lock), guard)) + FileWriteGuard::new(Arc::clone(&file_lock), guard) } // Tries to acquire a write lock for the specified file and returns a FileWriteGuard. - pub async fn try_write_lock(&self, path: &Path) -> io::Result { - let file_lock = self.get_or_create_lock(path).await.map_err(|_|str_to_io_error("Failed to acquire write lock"))?; + pub fn try_write_lock(&self, path: &Path) -> io::Result { + let file_lock = self.get_or_create_lock(path); let guard = file_lock.try_write(); match guard { // Clone the Arc to avoid moving `file_lock` out, as it is still borrowed by `guard` @@ -46,17 +46,17 @@ impl FileLockManager { // Helper function: retrieves or creates a lock for a file. - async fn get_or_create_lock(&self, path: &Path) -> io::Result>> { - let mut locks = self.locks.lock().await; + fn get_or_create_lock(&self, path: &Path) -> Arc> { + let mut locks = self.locks.lock(); if let Some(lock) = locks.get(path) { - return Ok(lock.clone()); + return lock.clone(); } let file_lock = Arc::new(RwLock::new(())); locks.insert(path.to_path_buf(), file_lock.clone()); drop(locks); - Ok(file_lock) + file_lock } } diff --git a/src/utils/file/mod.rs b/src/utils/file/mod.rs index 8afe663ac..c26f5ccfd 100644 --- a/src/utils/file/mod.rs +++ b/src/utils/file/mod.rs @@ -1,3 +1,4 @@ pub mod file_utils; pub mod multi_file_reader; -pub mod file_lock_manager; \ No newline at end of file +pub mod file_lock_manager; +pub mod config_reader; \ No newline at end of file diff --git a/src/utils/mod.rs b/src/utils/mod.rs index 5e8613121..d156014b8 100644 --- a/src/utils/mod.rs +++ b/src/utils/mod.rs @@ -1,6 +1,5 @@ pub mod string_utils; pub mod json_utils; -pub mod config_reader; pub mod default_utils; pub mod size_utils; pub mod sys_utils; diff --git a/src/utils/network/request.rs b/src/utils/network/request.rs index f23206c24..a057730c5 100644 --- a/src/utils/network/request.rs +++ b/src/utils/network/request.rs @@ -1,4 +1,4 @@ -use crate::Arc; +use std::sync::Arc; use std::collections::{HashMap, HashSet}; use std::fs; use std::fs::File; diff --git a/src/utils/network/xtream.rs b/src/utils/network/xtream.rs index 29ed33bec..6c0fa07c4 100644 --- a/src/utils/network/xtream.rs +++ b/src/utils/network/xtream.rs @@ -1,4 +1,4 @@ -use crate::Arc; +use std::sync::Arc; use crate::m3u_filter_error::{str_to_io_error, M3uFilterError}; use crate::model::config::{Config, ConfigInput, ConfigTarget}; use crate::model::playlist::{PlaylistEntry, PlaylistGroup, XtreamCluster, XtreamPlaylistItem}; @@ -69,31 +69,31 @@ where P: PlaylistEntry, { if cluster == XtreamCluster::Series { - if let Some(content) = xtream_repository::xtream_load_series_info(config, target.name.as_str(), pli.get_virtual_id()).await { + if let Some(content) = xtream_repository::xtream_load_series_info(config, target.name.as_str(), pli.get_virtual_id()) { // Deliver existing target content - return rewrite_xtream_series_info_content(config, target, pli, user, &content).await; + return rewrite_xtream_series_info_content(config, target, pli, user, &content); } // Check if the content has been resolved let resolve_series = target.options.as_ref().is_some_and(|opt| opt.xtream_resolve_series); if resolve_series { if let Some(provider_id) = pli.get_provider_id() { - if let Some(content) = xtream_get_input_info(config, input, provider_id, XtreamCluster::Series).await { - return xtream_repository::write_and_get_xtream_series_info(config, target, pli, user, &content).await; + if let Some(content) = xtream_get_input_info(config, input, provider_id, XtreamCluster::Series) { + return xtream_repository::write_and_get_xtream_series_info(config, target, pli, user, &content); } } } } else if cluster == XtreamCluster::Video { - if let Some(content) = xtream_repository::xtream_load_vod_info(config, target.name.as_str(), pli.get_virtual_id()).await { + if let Some(content) = xtream_repository::xtream_load_vod_info(config, target.name.as_str(), pli.get_virtual_id()) { // Deliver existing target content - return rewrite_xtream_vod_info_content(config, target, pli, user, &content).await; + return rewrite_xtream_vod_info_content(config, target, pli, user, &content); } // Check if the content has been resolved let resolve_vod = target.options.as_ref().is_some_and(|opt| opt.xtream_resolve_vod); if resolve_vod { if let Some(provider_id) = pli.get_provider_id() { - if let Some(content) = xtream_get_input_info(config, input, provider_id, XtreamCluster::Video).await { - return xtream_repository::write_and_get_xtream_vod_info(config, target, pli, user, &content).await; + if let Some(content) = xtream_get_input_info(config, input, provider_id, XtreamCluster::Video) { + return xtream_repository::write_and_get_xtream_vod_info(config, target, pli, user, &content); } } } @@ -102,8 +102,8 @@ where if let Ok(content) = get_xtream_stream_info_content(client, info_url, input).await { return match cluster { XtreamCluster::Live => Ok(content), - XtreamCluster::Video => xtream_repository::write_and_get_xtream_vod_info(config, target, pli, user, &content).await, - XtreamCluster::Series => xtream_repository::write_and_get_xtream_series_info(config, target, pli, user, &content).await, + XtreamCluster::Video => xtream_repository::write_and_get_xtream_vod_info(config, target, pli, user, &content), + XtreamCluster::Series => xtream_repository::write_and_get_xtream_series_info(config, target, pli, user, &content), }; }