diff --git a/Cargo.lock b/Cargo.lock index f2821dd7a..d8beaddde 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -137,9 +137,9 @@ dependencies = [ [[package]] name = "async-compression" -version = "0.4.20" +version = "0.4.21" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "310c9bcae737a48ef5cdee3174184e6d548b292739ede61a1f955ef76a738861" +checksum = "c0cf008e5e1a9e9e22a7d3c9a4992e21a350290069e36d8fb72304ed17e8f2d2" dependencies = [ "brotli", "flate2", @@ -301,9 +301,9 @@ dependencies = [ [[package]] name = "blake3" -version = "1.6.1" +version = "1.7.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "675f87afced0413c9bb02843499dbbd3882a237645883f71a2b59644a6d2f753" +checksum = "b17679a8d69b6d7fd9cd9801a536cec9fa5e5970b69f9d4747f70b39b031f5e7" dependencies = [ "arrayref", "arrayvec", @@ -587,9 +587,9 @@ dependencies = [ [[package]] name = "deranged" -version = "0.3.11" +version = "0.4.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "b42b6fa04a440b495c8b04d0e71b707c585f83cb9cb28cf8cd0d976c315e31b4" +checksum = "9c9e6a11ca8224451684bc0d7d5a7adbf8f2fd6887261a1cfc3c0432f9d4068e" dependencies = [ "powerfmt", ] @@ -916,14 +916,16 @@ dependencies = [ [[package]] name = "getrandom" -version = "0.3.1" +version = "0.3.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "43a49c392881ce6d5c3b8cb70f98717b7c07aabbdff06687b9030dbfbe2725f8" +checksum = "73fea8450eea4bac3940448fb7ae50d91f034f941199fcd9d909a5a07aa455f0" dependencies = [ "cfg-if", + "js-sys", "libc", - "wasi 0.13.3+wasi-0.2.2", - "windows-targets 0.52.6", + "r-efi", + "wasi 0.14.2+wasi-0.2.4", + "wasm-bindgen", ] [[package]] @@ -1421,9 +1423,9 @@ dependencies = [ [[package]] name = "libz-sys" -version = "1.1.21" +version = "1.1.22" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "df9b68e50e6e0b26f672573834882eb57759f6db9b3be2ea3c35c91188bb4eaa" +checksum = "8b70e7a7df205e92a1a4cd9aaae7898dac0aa555503cc0a649494d0d60e7651d" dependencies = [ "cc", "libc", @@ -1433,9 +1435,9 @@ dependencies = [ [[package]] name = "linux-raw-sys" -version = "0.9.2" +version = "0.9.3" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6db9c683daf087dc577b7506e9695b3d556a9f3849903fa28186283afd6809e9" +checksum = "fe7db12097d22ec582439daf8618b8fdd1a7bef6270e9af3b1ebcd30893cf413" [[package]] name = "litemap" @@ -1485,7 +1487,7 @@ dependencies = [ "pest", "pest_derive", "quick-xml", - "rand 0.9.0", + "rand", "regex", "reqwest", "rpassword", @@ -1503,7 +1505,6 @@ dependencies = [ "tower-http", "unidecode", "url", - "uuid", "vergen", "winapi", ] @@ -1627,9 +1628,9 @@ dependencies = [ [[package]] name = "once_cell" -version = "1.21.0" +version = "1.21.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "cde51589ab56b20a6f686b2c68f7a0bd6add753d697abf720d63f8db3ab7b1ad" +checksum = "d75b0bedcc4fe52caa0e03d9f1151a323e4aa5e2d78ba3580400cd3c9e2bc4bc" [[package]] name = "openssl" @@ -1892,11 +1893,12 @@ dependencies = [ [[package]] name = "quinn" -version = "0.11.6" +version = "0.11.7" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "62e96808277ec6f97351a2380e6c25114bc9e67037775464979f3037c92d05ef" +checksum = "c3bd15a6f2967aef83887dcb9fec0014580467e33720d073560cf015a5683012" dependencies = [ "bytes", + "cfg_aliases", "pin-project-lite", "quinn-proto", "quinn-udp", @@ -1906,17 +1908,18 @@ dependencies = [ "thiserror", "tokio", "tracing", + "web-time", ] [[package]] name = "quinn-proto" -version = "0.11.9" +version = "0.11.10" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "a2fe5ef3495d7d2e377ff17b1a8ce2ee2ec2a18cde8b6ad6619d65d0701c135d" +checksum = "b820744eb4dc9b57a3398183639c511b5a26d2ed702cedd3febaa1393caa22cc" dependencies = [ "bytes", - "getrandom 0.2.15", - "rand 0.8.5", + "getrandom 0.3.2", + "rand", "ring", "rustc-hash", "rustls", @@ -1952,15 +1955,10 @@ dependencies = [ ] [[package]] -name = "rand" -version = "0.8.5" +name = "r-efi" +version = "5.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "34af8d1a0e25924bc5b7c43c079c942339d8f0a8b57c39049bef581b46327404" -dependencies = [ - "libc", - "rand_chacha 0.3.1", - "rand_core 0.6.4", -] +checksum = "74765f6d916ee2faa39bc8e68e4f3ed8949b48cccdac59983d287a7cb71ce9c5" [[package]] name = "rand" @@ -1968,21 +1966,11 @@ 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.3", + "rand_chacha", + "rand_core", "zerocopy", ] -[[package]] -name = "rand_chacha" -version = "0.3.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e6c10a63a0fa32252be49d21e7709d4d4baf8d231c2dbce1eaa8141b9b127d88" -dependencies = [ - "ppv-lite86", - "rand_core 0.6.4", -] - [[package]] name = "rand_chacha" version = "0.9.0" @@ -1990,16 +1978,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d3022b5f1df60f26e1ffddd6c66e8aa15de382ae63b3a0c1bfc0e4d3e3f325cb" dependencies = [ "ppv-lite86", - "rand_core 0.9.3", -] - -[[package]] -name = "rand_core" -version = "0.6.4" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "ec0be4795e2f6a28069bec0b5ff3e2ac9bafc99e6a9a7dc3547996c5c816922c" -dependencies = [ - "getrandom 0.2.15", + "rand_core", ] [[package]] @@ -2008,7 +1987,7 @@ version = "0.9.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "99d9a13982dcf210057a8a78572b2217b667c3beacbf3a0d8b454f6f82837d38" dependencies = [ - "getrandom 0.3.1", + "getrandom 0.3.2", ] [[package]] @@ -2051,9 +2030,9 @@ checksum = "2b15c43186be67a4fd63bee50d0303afffcef381492ebe2c5d87f324e1b8815c" [[package]] name = "reqwest" -version = "0.12.14" +version = "0.12.15" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "989e327e510263980e231de548a33e63d34962d29ae61b467389a1a09627a254" +checksum = "d19c46a6fdd48bc4dab94b6103fccc55d34c67cc0ad04653aad4ea2a07cd7bbb" dependencies = [ "base64 0.22.1", "bytes", @@ -2174,9 +2153,9 @@ dependencies = [ [[package]] name = "rustix" -version = "1.0.2" +version = "1.0.3" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "f7178faa4b75a30e269c71e61c353ce2748cf3d76f0c44c393f4e60abf49b825" +checksum = "e56a18552996ac8d29ecc3b190b4fdbb2d91ca4ec396de7bbffaf43f3d637e96" dependencies = [ "bitflags 2.9.0", "errno", @@ -2187,9 +2166,9 @@ dependencies = [ [[package]] name = "rustls" -version = "0.23.23" +version = "0.23.25" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "47796c98c480fce5406ef69d1c76378375492c3b0a0de587be0c1d9feb12f395" +checksum = "822ee9188ac4ec04a2f0531e55d035fb2de73f18b41a63c70c2712503b6fb13c" dependencies = [ "once_cell", "ring", @@ -2219,9 +2198,9 @@ dependencies = [ [[package]] name = "rustls-webpki" -version = "0.102.8" +version = "0.103.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "64ca1bc8749bd4cf37b5ce386cc146580777b4e8572c7b97baf22c83f444bee9" +checksum = "0aa4eeac2588ffff23e9d7a7e9b3f971c5fb5b7ebc9452745e0c232c64f83b2f" dependencies = [ "ring", "rustls-pki-types", @@ -2491,13 +2470,12 @@ dependencies = [ [[package]] name = "tempfile" -version = "3.18.0" +version = "3.19.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "2c317e0a526ee6120d8dabad239c8dadca62b24b6f168914bbbc8e2fb1f0e567" +checksum = "488960f40a3fd53d72c2a29a58722561dee8afdd175bd88e3db4677d7b2ba600" dependencies = [ - "cfg-if", "fastrand 2.3.0", - "getrandom 0.3.1", + "getrandom 0.3.2", "once_cell", "rustix", "windows-sys 0.59.0", @@ -2525,9 +2503,9 @@ dependencies = [ [[package]] name = "time" -version = "0.3.39" +version = "0.3.40" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "dad298b01a40a23aac4580b67e3dbedb7cc8402f3592d7f49469de2ea4aecdd8" +checksum = "9d9c75b47bdff86fa3334a3db91356b8d7d86a9b839dab7d0bdc5c3d3a077618" dependencies = [ "deranged", "itoa", @@ -2542,15 +2520,15 @@ dependencies = [ [[package]] name = "time-core" -version = "0.1.3" +version = "0.1.4" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "765c97a5b985b7c11d7bc27fa927dc4fe6af3a6dfb021d28deb60d3bf51e76ef" +checksum = "c9e9a38711f559d9e3ce1cdb06dd7c5b8ea546bc90052da6d06bb76da74bb07c" [[package]] name = "time-macros" -version = "0.2.20" +version = "0.2.21" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e8093bc3e81c3bc5f7879de09619d06c9a5a5e45ca44dfeeb7225bae38005c5c" +checksum = "29aa485584182073ed57fd5004aa09c371f021325014694e432313345865fd04" dependencies = [ "num-conv", "time-core", @@ -2841,15 +2819,6 @@ version = "0.2.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "06abde3611657adf66d383f00b093d7faecc7fa57071cce2578660c9f1010821" -[[package]] -name = "uuid" -version = "1.16.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "458f7a779bf54acc9f347480ac654f68407d3aab21269a6e3c9f922acd9e2da9" -dependencies = [ - "getrandom 0.3.1", -] - [[package]] name = "vcpkg" version = "0.2.15" @@ -2915,9 +2884,9 @@ checksum = "9c8d87e72b64a3b4db28d11ce29237c246188f4f51057d65a7eab63b7987e423" [[package]] name = "wasi" -version = "0.13.3+wasi-0.2.2" +version = "0.14.2+wasi-0.2.4" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "26816d2e1a4a36a2940b96c5296ce403917633dff8f3440e9b236ed6f6bacad2" +checksum = "9683f9a5a998d873c0d21fcbe3c083009670149a8fab228644b8bd36b2c48cb3" dependencies = [ "wit-bindgen-rt", ] @@ -3068,9 +3037,9 @@ dependencies = [ [[package]] name = "windows-link" -version = "0.1.0" +version = "0.1.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6dccfd733ce2b1753b03b6d3c65edf020262ea35e20ccdf3e288043e6dd620e3" +checksum = "76840935b766e1b0a05c0066835fb9ec80071d4c09a16f6bd5f7e655e3c14c38" [[package]] name = "windows-registry" @@ -3085,9 +3054,9 @@ dependencies = [ [[package]] name = "windows-result" -version = "0.3.1" +version = "0.3.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "06374efe858fab7e4f881500e6e86ec8bc28f9462c47e5a9941a0142ad86b189" +checksum = "c64fd11a4fd95df68efcfee5f44a294fe71b8bc6a91993e2791938abcc712252" dependencies = [ "windows-link", ] @@ -3324,9 +3293,9 @@ dependencies = [ [[package]] name = "wit-bindgen-rt" -version = "0.33.0" +version = "0.39.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "3268f3d866458b787f390cf61f4bbb563b922d091359f9608842999eaee3943c" +checksum = "6f42320e61fe2cfd34354ecb597f86f413484a798ba44a8ca1165c58d42da6c1" dependencies = [ "bitflags 2.9.0", ] diff --git a/Cargo.toml b/Cargo.toml index 5806ba5a9..65a5a91c6 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -40,20 +40,19 @@ mime = "0.3" log = "0.4" env_logger = "0.11" rustelebot = "0.3" -bincode = { version = "2.0.0-rc.3", features = ["std", "serde"] } +bincode = { version = "2.0.1", features = ["std", "serde"] } rand = "0.9" rpassword = "7.3" flate2 = "1" -blake3 = "1.5" +blake3 = "1.7" bytes = "1.10" tokio-stream = { version = "0.1", features = ["sync"] } -tokio = { version = "1.43", features = ["rt-multi-thread", "parking_lot", "fs"] } +tokio = { version = "1.44", features = ["rt-multi-thread", "parking_lot", "fs"] } tokio-util = "0.7" paste = "1.0" -tempfile = "3.16" +tempfile = "3.19" ruzstd = "0" filetime = "0.2" -uuid = { version = "1", features = ["v4"] } #[cfg(target_os = "macos")] libc = "0" #[cfg(target_os = "windows")] diff --git a/bin/build_resources.sh b/bin/build_resources.sh index 889100321..90f87f684 100755 --- a/bin/build_resources.sh +++ b/bin/build_resources.sh @@ -21,7 +21,7 @@ while getopts "fh" opt; do esac done -declare -a resources=("channel_unavailable" "user_connections_exhausted") +declare -a resources=("channel_unavailable" "user_connections_exhausted" "provider_connections_exhausted") for resource in "${resources[@]}"; do if [ "$flag_force" = false ]; then diff --git a/resources/provider_connections_exhausted.jpg b/resources/provider_connections_exhausted.jpg new file mode 100644 index 000000000..fde4b9228 Binary files /dev/null and b/resources/provider_connections_exhausted.jpg differ diff --git a/src/api/api_utils.rs b/src/api/api_utils.rs index 69ea68b12..3e50c3617 100644 --- a/src/api/api_utils.rs +++ b/src/api/api_utils.rs @@ -68,6 +68,8 @@ macro_rules! try_result_bad_request { pub use try_option_bad_request; pub use try_result_bad_request; +use crate::api::model::active_provider_manager::ProviderConfig; +use crate::api::model::streams::provider_stream::{create_provider_connections_exhausted_stream, ProviderStreamResponse}; use crate::auth::authenticator::Claims; pub async fn serve_file(file_path: &Path, mime_type: mime::Mime) -> impl axum::response::IntoResponse + Send { @@ -155,78 +157,141 @@ fn get_stream_options(app_state: &AppState) -> StreamOptions { // content_length // } -pub async fn stream_response(app_state: &AppState, stream_url: &str, +fn get_stream_alternative_url(stream_url: &str, input: &ConfigInput, alias_input: &ProviderConfig) -> String { + let Some(input_user_info) = input.get_user_info() else { return stream_url.to_owned() }; + let Some(alt_input_user_info) = alias_input.get_user_info() else { return stream_url.to_owned() }; + + let modified = stream_url.replace(&input_user_info.base_url, &alt_input_user_info.base_url); + let modified = modified.replace(&input_user_info.username, &alt_input_user_info.username); + let modified = modified.replace(&input_user_info.password, &alt_input_user_info.password); + modified +} + +/** +* If successfully a provider connection is used, do not forget to release if unsuccessfully +*/ +fn get_stream_response_params(app_state: &AppState, stream_url: &str, input_opt: Option<&ConfigInput>) + -> (Option, Option>, Option, Option) { + let (custom_stream, input_headers, input_name, request_url) = if let Some(input) = input_opt { + let (stream_response, alias_input_name, request_url) = match app_state.active_provider.acquire_connection(&input.name) { + None => { + let stream = create_provider_connections_exhausted_stream(&app_state.config, &[]); + (Some(stream), None, None) + }, + Some(alias_input) => { + if alias_input.id != input.id { + (None, Some(alias_input.name.to_string()), Some(get_stream_alternative_url(stream_url, input, &alias_input))) + } else { + (None, Some(input.name.to_string()), Some(stream_url.to_string())) + } + } + }; + (stream_response, Some(input.headers.clone()), alias_input_name, request_url) + } else { + (None, None, None, None) + }; + (custom_stream, input_headers, input_name, request_url) +} + +async fn create_stream_response(app_state: &AppState, stream_options: &StreamOptions, stream_url: &str, + req_headers: &HeaderMap, input_opt: Option<&ConfigInput>, + item_type: PlaylistItemType, share_stream : bool) -> (ProviderStreamResponse, Option) { + let (provider_stream, input_headers, stream_input_name, request_url) = get_stream_response_params(app_state, stream_url, input_opt); + if let Some(provider_stream_response) = provider_stream{ + return (provider_stream_response, None); + } + + let parsed_url = request_url.map_or_else(|| Url::parse(stream_url), |u| Url::parse(&u)); + + let (stream, stream_info) = if let Ok(url) = parsed_url { + if stream_options.pipe_provider_stream { + provider_stream::get_provider_pipe_stream(app_state, &url, req_headers, input_headers, item_type, &stream_options).await + } else { + let buffer_stream_options = BufferStreamOptions::new(item_type, share_stream, &stream_options); + provider_stream::get_provider_reconnect_buffered_stream(app_state, &url, req_headers, input_headers, buffer_stream_options).await + } + } else { + (None, None) + }; + + // if we have no stream we should release the provider + if stream.is_none() { + if let Some(alt_input_name) = &stream_input_name { + app_state.active_provider.release_connection(alt_input_name); + } + } + + ((stream, stream_info), stream_input_name.or(input_opt.map(|i| i.name.to_string()))) +} + +pub async fn stream_response(app_state: &AppState, + stream_url: &str, req_headers: &HeaderMap, input: Option<&ConfigInput>, - item_type: PlaylistItemType, target: &ConfigTarget, + item_type: PlaylistItemType, + target: &ConfigTarget, user: &ProxyUserCredentials) -> impl axum::response::IntoResponse + Send { if log_enabled!(log::Level::Trace) { trace!("Try to open stream {}", sanitize_sensitive_info(stream_url)); } let share_stream = is_stream_share_enabled(item_type, target); if share_stream { - if let Some(value) = shared_stream_response(app_state, stream_url, user, input).await { + if let Some(value) = shared_stream_response(app_state, stream_url, user).await { return value.into_response(); } } - let stream_options = get_stream_options(app_state); get_stream_options(app_state); - - if let Ok(url) = Url::parse(stream_url) { - let event_manager = Arc::clone(&app_state.event_manager); - let (stream_opt, provider_response) = if stream_options.pipe_provider_stream { - provider_stream::get_provider_pipe_stream(&app_state.config, &app_state.http_client, &url, req_headers, input, item_type, &stream_options).await - } else { - let buffer_stream_options = BufferStreamOptions::new(item_type,share_stream, &stream_options); - provider_stream::get_provider_reconnect_buffered_stream(&app_state.config, &app_state.http_client, &url, req_headers, input, buffer_stream_options).await - }; - if let Some(stream) = stream_opt { - // let content_length = get_stream_content_length(provider_response.as_ref()); - let stream = ActiveClientStream::new(stream, event_manager, &user.username, input.map(|c| c.name.clone())).await; - 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_headers, stream_options.buffer_size).await; - if let Some(broadcast_stream) = SharedStreamManager::subscribe_shared_stream(app_state, stream_url).await { - let (status_code, header_map) = get_stream_response_with_headers(provider_response, stream_url); - let mut response = axum::response::Response::builder() - .status(status_code); - for (key, value) in &header_map { - response = response.header(key, value); - } - response.body(axum::body::Body::from_stream(broadcast_stream)).unwrap().into_response() - // if content_length > 0 { - // response_builder.body(SizedStream::new(content_length, broadcast_stream)) } - // else { - // response_builder.body(BodyStream::new(broadcast_stream)) - // } - } else { - axum::http::StatusCode::BAD_REQUEST.into_response() - } - } else { + let stream_options = get_stream_options(app_state); + let event_manager = Arc::clone(&app_state.event_manager); + let ((stream_opt, provider_response), stream_input_name) = create_stream_response(app_state, &stream_options, &stream_url, req_headers, input, item_type, share_stream).await; + if let Some(stream) = stream_opt { + // let content_length = get_stream_content_length(provider_response.as_ref()); + let stream = ActiveClientStream::new(stream, event_manager, &user.username, stream_input_name).await; + let stream_resp = if share_stream { + // Shared Stream response + let shared_headers = provider_response.as_ref().map_or_else(Vec::new, |(h, _)| h.clone()); + SharedStreamManager::subscribe(app_state, stream_url, stream, shared_headers, stream_options.buffer_size).await; + if let Some(broadcast_stream) = SharedStreamManager::subscribe_shared_stream(app_state, stream_url).await { let (status_code, header_map) = get_stream_response_with_headers(provider_response, stream_url); let mut response = axum::response::Response::builder() .status(status_code); for (key, value) in &header_map { response = response.header(key, value); } - response.body(axum::body::Body::from_stream(stream)).unwrap().into_response() + response.body(axum::body::Body::from_stream(broadcast_stream)).unwrap().into_response() + // if content_length > 0 { + // response_builder.body(SizedStream::new(content_length, broadcast_stream)) } + // else { + // response_builder.body(BodyStream::new(broadcast_stream)) + // } + } else { + axum::http::StatusCode::BAD_REQUEST.into_response() + } + } else { + let (status_code, header_map) = get_stream_response_with_headers(provider_response, stream_url); + let mut response = axum::response::Response::builder() + .status(status_code); + for (key, value) in &header_map { + response = response.header(key, value); + } + response.body(axum::body::Body::from_stream(stream)).unwrap().into_response() - // if content_length > 0 { response_builder.body(SizedStream::new(content_length, stream)) } else { response_builder.streaming(stream) } - }; + // if content_length > 0 { response_builder.body(SizedStream::new(content_length, stream)) } else { response_builder.streaming(stream) } + }; - return stream_resp.into_response(); - } + return stream_resp.into_response(); } + error!("Cant open stream {}", sanitize_sensitive_info(stream_url)); axum::http::StatusCode::BAD_REQUEST.into_response() } -async fn shared_stream_response(app_state: &AppState, stream_url: &str, user: &ProxyUserCredentials, input: Option<&ConfigInput>,) -> Option { +async fn shared_stream_response(app_state: &AppState, stream_url: &str, user: &ProxyUserCredentials) -> Option { if let Some(stream) = SharedStreamManager::subscribe_shared_stream(app_state, stream_url).await { debug_if_enabled!("Using shared channel {}", sanitize_sensitive_info(stream_url)); if let Some(headers) = app_state.shared_stream_manager.get_shared_state_headers(stream_url).await { let (status_code, header_map) = get_stream_response_with_headers(Some((headers.clone(), StatusCode::OK)), stream_url); let event_manager = Arc::clone(&app_state.event_manager); - let stream = ActiveClientStream::new(stream, event_manager, &user.username, input.map(|c| c.name.clone())).await.boxed(); + let stream = ActiveClientStream::new(stream, event_manager, &user.username, None).await.boxed(); let mut response = axum::response::Response::builder() .status(status_code); for (key, value) in &header_map { @@ -239,7 +304,7 @@ async fn shared_stream_response(app_state: &AppState, stream_url: &str, user: &P } pub fn is_stream_share_enabled(item_type: PlaylistItemType, target: &ConfigTarget) -> bool { - (item_type == PlaylistItemType::Live || item_type == PlaylistItemType::LiveHls) && target.options.as_ref().is_some_and(|opt| opt.share_live_streams) + (item_type == PlaylistItemType::Live /* || item_type == PlaylistItemType::LiveHls */) && target.options.as_ref().is_some_and(|opt| opt.share_live_streams) } pub type HeaderFilter = Option bool + Send>>; diff --git a/src/api/endpoints/hdhomerun_api.rs b/src/api/endpoints/hdhomerun_api.rs index 711040766..1918e797b 100644 --- a/src/api/endpoints/hdhomerun_api.rs +++ b/src/api/endpoints/hdhomerun_api.rs @@ -93,9 +93,9 @@ where match channels { Some(chans) => { let mapped = chans.map(move |(item, has_next)| { - let input = cfg.get_input_by_name(&item.input_name); - let (live_stream_use_prefix, live_stream_without_extension) = input.map_or((true, false), |i| i.options.as_ref() - .map_or((true, false), |o| (o.xtream_live_stream_use_prefix, o.xtream_live_stream_without_extension))); + let input_options = cfg.get_input_options_by_name(&item.input_name); + let (live_stream_use_prefix, live_stream_without_extension) = input_options.as_ref() + .map_or((true, false), |o| (o.xtream_live_stream_use_prefix, o.xtream_live_stream_without_extension)); let container_extension = item.get_additional_property("container_extension").map(|v| get_string_from_serde_value(&v).unwrap_or_default()); let stream_url = match &base_url { None => item.url.to_string(), diff --git a/src/api/endpoints/hls_api.rs b/src/api/endpoints/hls_api.rs index 30d8838e4..2ae553367 100644 --- a/src/api/endpoints/hls_api.rs +++ b/src/api/endpoints/hls_api.rs @@ -15,7 +15,7 @@ use std::sync::Arc; #[derive(Deserialize)] struct HlsApiPathParams { - token: String, + token: u32, username: String, password: String, stream_id: u32, @@ -32,7 +32,8 @@ pub(in crate::api) async fn handle_hls_stream_request(app_state: &Arc, let server_info = app_state.config.get_user_server_info(user).await; match request::download_text_content(Arc::clone(&app_state.http_client), input, &url, None).await { Ok(content) => { - let (hls_entry, hls_content) = rewrite_hls(&server_info.get_base_url(), &content, hls_url, virtual_id, user, &target_type, &input.name); + let hls_token = app_state.hls_cache.new_token(); + let (hls_entry, hls_content) = rewrite_hls(&server_info.get_base_url(), &content, hls_url, virtual_id, hls_token, user, &target_type, input.id); app_state.hls_cache.add_entry(hls_entry).await; axum::response::Response::builder() .status(axum::http::StatusCode::OK) @@ -62,12 +63,11 @@ async fn hls_api_stream( if user.connections_exhausted(&app_state).await { return create_custom_video_stream_response(&app_state.config, &CustomVideoStreamType::UserConnectionsExhausted).into_response(); } - - let Some(hls_entry) = app_state.hls_cache.get_entry(¶ms.token).await else { return axum::http::StatusCode::BAD_REQUEST.into_response(); }; + let Some(hls_entry) = app_state.hls_cache.get_entry(params.token).await else { return axum::http::StatusCode::BAD_REQUEST.into_response(); }; let Some(hls_url) = hls_entry.get_chunk_url(params.chunk) else { return axum::http::StatusCode::BAD_REQUEST.into_response(); }; let target_name = &target.name; let virtual_id = params.stream_id; - let input = try_option_bad_request!(app_state.config.get_input_by_name(&hls_entry.input_name), true, format!("Cant find input for target {target_name}, context {}, stream_id {virtual_id}", XtreamCluster::Live)); + let input = try_option_bad_request!(app_state.config.get_input_by_id(hls_entry.input_id), true, format!("Cant find input for target {target_name}, context {}, stream_id {virtual_id}", XtreamCluster::Live)); if is_hls_url(hls_url) { return handle_hls_stream_request(&app_state, &user, hls_url, virtual_id, input, hls_entry.target_type.clone()).await.into_response(); @@ -80,7 +80,7 @@ async fn hls_api_stream( // let pli = try_result_bad_request!(m3u_repository::m3u_get_item_for_stream_id(virtual_id, &app_state.config, target).await, true, format!("Failed to read xtream item for stream id {}", virtual_id)); // (pli.url, pli.input_name) // }; - stream_response(&app_state, hls_url, &req_headers, Some(input), PlaylistItemType::Live, target, &user).await.into_response() + stream_response(&app_state, hls_url, &req_headers, Some(input), PlaylistItemType::LiveHls, target, &user).await.into_response() } pub fn hls_api_register() -> axum::Router> { diff --git a/src/api/endpoints/m3u_api.rs b/src/api/endpoints/m3u_api.rs index 13dc0286e..2a11b2710 100644 --- a/src/api/endpoints/m3u_api.rs +++ b/src/api/endpoints/m3u_api.rs @@ -7,7 +7,7 @@ use crate::model::config::TargetType; use crate::model::playlist::{FieldGetAccessor, PlaylistItemType, XtreamCluster}; use crate::repository::m3u_playlist_iterator::{M3U_RESOURCE_PATH, M3U_STREAM_PATH}; use crate::repository::m3u_repository::{m3u_get_item_for_stream_id, m3u_load_rewrite_playlist}; -use crate::utils::network::request::{replace_url_extension, sanitize_sensitive_info, HLS_EXT}; +use crate::utils::network::request::{replace_url_extension, sanitize_sensitive_info, DASH_EXT, HLS_EXT}; use crate::utils::debug_if_enabled; use axum::response::IntoResponse; use bytes::Bytes; @@ -91,13 +91,16 @@ async fn m3u_api_stream( let input = app_state.config.get_input_by_name(m3u_item.input_name.as_str()); let is_hls_request = m3u_item.item_type == PlaylistItemType::LiveHls || stream_ext.as_deref() == Some(HLS_EXT); + let is_dash_request = !is_hls_request && m3u_item.item_type == PlaylistItemType::LiveDash || stream_ext.as_deref() == Some(DASH_EXT); - if user.proxy == ProxyType::Redirect { - let redirect_url = if is_hls_request { &replace_url_extension(&m3u_item.url, "m3u8") } else { &m3u_item.url }; + if user.proxy == ProxyType::Redirect || is_dash_request { + let redirect_url = if is_hls_request { &replace_url_extension(&m3u_item.url, HLS_EXT) } else { &m3u_item.url }; + let redirect_url = if is_dash_request { &replace_url_extension(redirect_url, DASH_EXT) } else { redirect_url }; // TODO alias processing debug_if_enabled!("Redirecting m3u stream request to {}", sanitize_sensitive_info(redirect_url)); return redirect(redirect_url.as_str()).into_response(); } + // Reverse proxy mode if is_hls_request { let target_name = &target.name; diff --git a/src/api/endpoints/xtream_api.rs b/src/api/endpoints/xtream_api.rs index 4ecd45a70..0ec8ea324 100644 --- a/src/api/endpoints/xtream_api.rs +++ b/src/api/endpoints/xtream_api.rs @@ -35,7 +35,7 @@ use crate::repository::{user_repository, xtream_repository}; use crate::repository::xtream_repository::{TAG_EPISODES, TAG_INFO_DATA, TAG_SEASONS_DATA}; use crate::utils::hash_utils::generate_playlist_uuid; use crate::utils::json_utils::get_u32_from_serde_value; -use crate::utils::network::request::{extract_extension_from_url, sanitize_sensitive_info, HLS_EXT}; +use crate::utils::network::request::{extract_extension_from_url, replace_url_extension, sanitize_sensitive_info, DASH_EXT, HLS_EXT}; use crate::utils::network::xtream::{create_vod_info_from_item, ACTION_GET_LIVE_CATEGORIES, ACTION_GET_LIVE_STREAMS, ACTION_GET_SERIES, ACTION_GET_SERIES_CATEGORIES, ACTION_GET_SERIES_INFO, ACTION_GET_VOD_CATEGORIES, ACTION_GET_VOD_INFO, ACTION_GET_VOD_STREAMS}; use crate::utils::json_utils; use crate::utils::debug_if_enabled; @@ -108,7 +108,7 @@ pub fn serve_query(file_path: &Path, filter: &HashMap<&str, HashSet>) -> } fn get_xtream_player_api_stream_url(input: &ConfigInput, context: &XtreamApiStreamContext, action_path: &str, fallback_url: &str) -> Option { - if let Some(user_info) = input.get_user_info() { + if let Some(input_user_info) = input.get_user_info() { let ctx = match context { XtreamApiStreamContext::LiveAlt | XtreamApiStreamContext::Live => { @@ -121,10 +121,10 @@ fn get_xtream_player_api_stream_url(input: &ConfigInput, context: &XtreamApiStre }; let ctx_path = if ctx.is_empty() { String::new() } else { format!("{ctx}/") }; Some(format!("{}/{}{}/{}/{}", - &user_info.base_url, + &input_user_info.base_url, ctx_path, - &user_info.username, - &user_info.password, + &input_user_info.username, + &input_user_info.password, action_path )) } else if !fallback_url.is_empty() { @@ -177,10 +177,15 @@ async fn xtream_player_api_stream( } // if pli.item_type == PlaylistItemType::LiveHls { - // let redirect_url = &replace_extension(&pli.url, "m3u8"); + // let redirect_url = &replace_extension(&pli.url, HLS_EXT); // debug_if_enabled!("Redirecting stream request to {}", sanitize_sensitive_info(redirect_url)); // return redirect(redirect_url).into_response(); // } + if pli.item_type == PlaylistItemType::LiveDash { + let redirect_url = &replace_url_extension(&pli.url, DASH_EXT); + debug_if_enabled!("Redirecting stream request to {}", sanitize_sensitive_info(redirect_url)); + return redirect(redirect_url).into_response(); + } debug_if_enabled!("Redirecting stream request to {}", sanitize_sensitive_info(&pli.url)); return redirect(&pli.url).into_response(); diff --git a/src/api/model/active_provider_manager.rs b/src/api/model/active_provider_manager.rs index 3a8d22905..d76742280 100644 --- a/src/api/model/active_provider_manager.rs +++ b/src/api/model/active_provider_manager.rs @@ -1,5 +1,5 @@ use std::cell::RefCell; -use crate::model::config::{Config, ConfigInput, ConfigInputAlias, InputType}; +use crate::model::config::{Config, ConfigInput, ConfigInputAlias, InputType, InputUserInfo}; use std::collections::HashMap; use std::sync::atomic::{AtomicU16, AtomicUsize, Ordering}; @@ -52,6 +52,10 @@ impl ProviderConfig { } } + pub fn get_user_info(&self) -> Option { + InputUserInfo::new(self.input_type.clone(), self.username.as_deref(), self.password.as_deref(), &self.url) + } + #[inline] pub fn is_exhausted(&self) -> bool { self.max_connections > 0 && self.current_connections.load(Ordering::SeqCst) >= self.max_connections @@ -62,9 +66,9 @@ impl ProviderConfig { // !self.is_exhausted() // } - pub fn try_allocate(&self, force: bool) -> bool { + pub fn try_allocate(&self) -> bool { let connections = self.current_connections.load(Ordering::SeqCst); - if force || self.max_connections == 0 || connections < self.max_connections { + if self.max_connections == 0 || connections < self.max_connections { self.current_connections.fetch_add(1, Ordering::SeqCst); return true; } @@ -94,10 +98,10 @@ enum ProviderLineup { } impl ProviderLineup { - fn acquire(&self, force: bool) -> Option<&ProviderConfig> { + fn acquire(&self) -> Option<&ProviderConfig> { match self { - ProviderLineup::Single(lineup) => lineup.acquire(force), - ProviderLineup::Multi(lineup) => lineup.acquire(force), + ProviderLineup::Single(lineup) => lineup.acquire(), + ProviderLineup::Multi(lineup) => lineup.acquire(), } } @@ -122,8 +126,8 @@ impl SingleProviderLineup { } } - fn acquire(&self, force: bool) -> Option<&ProviderConfig> { - if self.provider.try_allocate(force) { + fn acquire(&self) -> Option<&ProviderConfig> { + if self.provider.try_allocate() { Some(&self.provider) } else { None @@ -235,7 +239,7 @@ impl MultiProviderLineup { fn acquire_next_provider_from_group(priority_group: &ProviderPriorityGroup) -> Option<&ProviderConfig> { match priority_group { ProviderPriorityGroup::SingleProviderGroup(p) => { - if p.try_allocate(false) { + if p.try_allocate() { return Some(p); } } @@ -245,7 +249,7 @@ impl MultiProviderLineup { for _ in 0..provider_count { let p = pg.get(idx).unwrap(); idx = (idx + 1) % provider_count; - if p.try_allocate(false) { + if p.try_allocate() { index.store(idx, Ordering::SeqCst); return Some(p); } @@ -258,9 +262,6 @@ impl MultiProviderLineup { /// Attempts to acquire a provider from the lineup based on priority and availability. /// - /// # Parameters - /// - `force`: A boolean flag indicating whether to force allocation even if all providers are exhausted. - /// /// # Returns /// - `Some(&ProviderConfig)`: A reference to the acquired provider if allocation was successful. /// - `None`: If no providers are available and `force` is `false`. @@ -286,7 +287,7 @@ impl MultiProviderLineup { /// println!("No available providers."); /// } /// ``` - fn acquire(&self, force: bool) -> Option<&ProviderConfig> { + fn acquire(&self) -> Option<&ProviderConfig> { let mut main_idx = self.index.load(Ordering::SeqCst); let provider_count = self.providers.len(); @@ -301,21 +302,17 @@ impl MultiProviderLineup { } } - if force { - let provider = &self.providers[main_idx]; - self.index.store((main_idx + 1) % provider_count, Ordering::SeqCst); + let provider = &self.providers[main_idx]; + self.index.store((main_idx + 1) % provider_count, Ordering::SeqCst); - return match provider { - ProviderPriorityGroup::SingleProviderGroup(p) => Some(p), - ProviderPriorityGroup::MultiProviderGroup(gindex, group) => { - let idx = gindex.load(Ordering::SeqCst); - gindex.store((idx + 1) % group.len(), Ordering::SeqCst); - group.get(idx) - } - }; - } - - None + return match provider { + ProviderPriorityGroup::SingleProviderGroup(p) => Some(p), + ProviderPriorityGroup::MultiProviderGroup(gindex, group) => { + let idx = gindex.load(Ordering::SeqCst); + gindex.store((idx + 1) % group.len(), Ordering::SeqCst); + group.get(idx) + } + }; } @@ -343,15 +340,12 @@ impl MultiProviderLineup { } pub struct ActiveProviderManager { - user_access_control: bool, providers: Vec, } impl ActiveProviderManager { pub fn new(cfg: &Config) -> Self { - let user_access_control = cfg.user_access_control; let mut this = Self { - user_access_control, providers: Vec::new(), }; for source in &cfg.sources { @@ -403,10 +397,11 @@ impl ActiveProviderManager { pub fn acquire_connection(&self, input_name: &str) -> Option<&ProviderConfig> { match self.get_provider_config(input_name) { None => None, - Some((lineup, _config)) => lineup.acquire(self.user_access_control) + Some((lineup, _config)) => lineup.acquire() } } + // we need the provider_name to exactly release this provider pub fn release_connection(&self, provider_name: &str) { if let Some((lineup, _config)) = self.get_provider_config(provider_name) { lineup.release(provider_name); diff --git a/src/api/model/event_manager.rs b/src/api/model/event_manager.rs index d2e7a5a04..9d099412d 100644 --- a/src/api/model/event_manager.rs +++ b/src/api/model/event_manager.rs @@ -31,21 +31,18 @@ impl EventManager { pub async fn fire(&self, event: Event) { match event { - Event::StreamConnect((username, input_name)) => { + Event::StreamConnect((username, _input_name)) => { let (client_count, connection_count) = self.active_user.add_connection(&username).await; if self.log_active_clients { info!("Active clients: {client_count}, active connections {connection_count}"); } - if let Some(input) = input_name { - // TODO this is the wrong place, move it later to the right place - self.active_provider.acquire_connection(&input); - } } Event::StreamDisconnect((username, input_name)) => { let (client_count, connection_count) = self.active_user.remove_connection(&username).await; if self.log_active_clients { info!("Active clients: {client_count}, active connections {connection_count}"); } + if let Some(input) = input_name { self.active_provider.release_connection(&input); } diff --git a/src/api/model/hls_cache.rs b/src/api/model/hls_cache.rs index ca038e514..90e13b5d2 100644 --- a/src/api/model/hls_cache.rs +++ b/src/api/model/hls_cache.rs @@ -2,6 +2,7 @@ use log::error; use std::collections::HashMap; use std::str::FromStr; use std::sync::Arc; +use std::sync::atomic::AtomicU32; use std::time::Duration; use chrono::Local; use cron::Schedule; @@ -11,6 +12,7 @@ use crate::utils::sys_utils::exit; use crate::model::hls::HlsEntry; const EXPIRE_DURATION: u64 = 600; // 10 minutes +const TOKEN_MAX: u32 = u32::MAX - u16::MAX as u32; fn start_garbage_collector(cache: &Arc) { let cache_clone = cache.clone(); @@ -32,13 +34,15 @@ fn start_garbage_collector(cache: &Arc) { } pub struct HlsCache { - pub entries: RwLock>, + pub entries: RwLock>, + counter: AtomicU32, } impl HlsCache { pub fn garbage_collected() -> Arc { let cache = Arc::new(Self { entries: RwLock::new(HashMap::new()), + counter: AtomicU32::new(1), }); start_garbage_collector(&cache); @@ -46,11 +50,11 @@ impl HlsCache { } pub async fn add_entry(&self, entry: HlsEntry) { - self.entries.write().await.insert(entry.token.to_string(), entry); + self.entries.write().await.insert(entry.token, entry); } - pub async fn get_entry(&self, token: &str) -> Option{ - self.entries.read().await.get(token).cloned() + pub async fn get_entry(&self, token: u32) -> Option{ + self.entries.read().await.get(&token).cloned() } pub async fn gc(&self) { @@ -58,4 +62,13 @@ impl HlsCache { // Remove all expired elements self.entries.write().await.retain(|_, entry| entry.ts > threshold); } + + pub fn new_token(&self) -> u32 { + let token = self.counter.fetch_add(1, std::sync::atomic::Ordering::SeqCst); + if token > TOKEN_MAX { + self.counter.store(1, std::sync::atomic::Ordering::SeqCst); + return 1; + } + return token; + } } \ No newline at end of file diff --git a/src/api/model/mod.rs b/src/api/model/mod.rs index 51439073d..922490c82 100644 --- a/src/api/model/mod.rs +++ b/src/api/model/mod.rs @@ -9,4 +9,4 @@ pub(crate) mod streams; pub(in crate::api) mod active_user_manager; pub(in crate::api) mod active_provider_manager; pub(in crate::api) mod event_manager; -pub(in crate::api) mod hls_cache; +pub(in crate::api) mod hls_cache; \ No newline at end of file diff --git a/src/api/model/streams/provider_stream.rs b/src/api/model/streams/provider_stream.rs index e0d102bb9..bcd8aedd5 100644 --- a/src/api/model/streams/provider_stream.rs +++ b/src/api/model/streams/provider_stream.rs @@ -1,9 +1,10 @@ +use std::collections::HashMap; use crate::api::api_utils::{get_headers_from_request, HeaderFilter, StreamOptions}; use crate::api::model::model_utils::get_response_headers; use crate::api::model::stream_error::StreamError; use crate::api::model::streams::custom_video_stream::CustomVideoStream; use crate::api::model::streams::provider_stream_factory::{create_provider_stream, BufferStreamOptions}; -use crate::model::config::{Config, ConfigInput}; +use crate::model::config::{Config}; use crate::model::playlist::PlaylistItemType; use crate::utils::debug_if_enabled; use crate::utils::network::request::{get_request_headers, sanitize_sensitive_info}; @@ -17,10 +18,11 @@ use std::time::Duration; use axum::http::HeaderMap; use axum::response::IntoResponse; use url::Url; +use crate::api::model::app_state::AppState; type BoxedProviderStream = BoxStream<'static, Result>; type ProviderStreamHeader = Vec<(String, String)>; -type ProviderStreamResponse = (Option, Option<(ProviderStreamHeader, StatusCode)>); +pub type ProviderStreamResponse = (Option, Option<(ProviderStreamHeader, StatusCode)>); pub enum CustomVideoStreamType { ChannelUnavailable, @@ -28,33 +30,33 @@ pub enum CustomVideoStreamType { ProviderConnectionsExhausted, } -fn create_video_stream(video: Option<&Arc>>, headers: &[(String, String)], log_message: &str) -> Option<(BoxedProviderStream, (ProviderStreamHeader, StatusCode))> { +fn create_video_stream(video: Option<&Arc>>, headers: &[(String, String)], log_message: &str) -> ProviderStreamResponse { if let Some(video) = video { debug!("{}", log_message); let mut response_headers: Vec<(String, String)> = headers.iter() .filter(|(key, _)| !(key.eq("content-type") || key.eq("content-length") || key.contains("range"))) .map(|(key, value)| (key.to_string(), value.to_string())).collect(); response_headers.push(("content-type".to_string(), "video/mp2t".to_string())); - Some((Box::pin(CustomVideoStream::new(Arc::clone(video))), (response_headers, StatusCode::OK))) + (Some(Box::pin(CustomVideoStream::new(Arc::clone(video)))), Some((response_headers, StatusCode::OK))) } else { - None + (None, None) } } -pub fn create_channel_unavailable_stream(cfg: &Config, headers: &[(String, String)], status: StatusCode) -> Option<(BoxedProviderStream, (ProviderStreamHeader, StatusCode))> { +pub fn create_channel_unavailable_stream(cfg: &Config, headers: &[(String, String)], status: StatusCode) -> ProviderStreamResponse { create_video_stream(cfg.t_channel_unavailable_video.as_ref(), headers, &format!("Streaming response channel unavailable for status {status}")) } -pub fn create_user_connections_exhausted_stream(cfg: &Config, headers: &[(String, String)]) -> Option<(BoxedProviderStream, (ProviderStreamHeader, StatusCode))> { +pub fn create_user_connections_exhausted_stream(cfg: &Config, headers: &[(String, String)]) -> ProviderStreamResponse { create_video_stream(cfg.t_user_connections_exhausted_video.as_ref(), headers, "Streaming response user connections exhausted") } -pub fn create_provider_connections_exhausted_stream(cfg: &Config, headers: &[(String, String)]) -> Option<(BoxedProviderStream, (ProviderStreamHeader, StatusCode))> { +pub fn create_provider_connections_exhausted_stream(cfg: &Config, headers: &[(String, String)]) -> ProviderStreamResponse { create_video_stream(cfg.t_provider_connections_exhausted_video.as_ref(), headers, "Streaming response provider connections exhausted") } pub fn create_custom_video_stream_response(config: &Config, video_response: &CustomVideoStreamType) -> impl axum::response::IntoResponse + Send { - if let Some((stream, (headers, status_code))) = match video_response { + if let (Some(stream), Some((headers, status_code))) = match video_response { CustomVideoStreamType::ChannelUnavailable => create_channel_unavailable_stream(config, &[], StatusCode::BAD_REQUEST), CustomVideoStreamType::UserConnectionsExhausted => create_user_connections_exhausted_stream(config, &[]), CustomVideoStreamType::ProviderConnectionsExhausted => create_provider_connections_exhausted_stream(config, &[]), @@ -68,33 +70,29 @@ pub fn create_custom_video_stream_response(config: &Config, video_response: &Cus } axum::http::StatusCode::FORBIDDEN.into_response() } - - pub fn get_header_filter_for_item_type(item_type: PlaylistItemType) -> HeaderFilter { match item_type { - PlaylistItemType::Live | PlaylistItemType::LiveUnknown | PlaylistItemType::LiveHls => { + PlaylistItemType::Live | PlaylistItemType::LiveHls | PlaylistItemType::LiveDash | PlaylistItemType::LiveUnknown => { Some(Box::new(|key| key != "accept-ranges" && key != "range" && key != "content-range")) } _ => None, } } -pub async fn get_provider_pipe_stream(cfg: &Config, - http_client: &Arc, +pub async fn get_provider_pipe_stream(app_state: &AppState, stream_url: &Url, req_headers: &HeaderMap, - input: Option<&ConfigInput>, + input_headers: Option>, item_type: PlaylistItemType, stream_options: &StreamOptions) -> ProviderStreamResponse { let filter_header = get_header_filter_for_item_type(item_type); let req_headers = get_headers_from_request(req_headers, &filter_header); debug_if_enabled!("Stream requested with headers: {:?}", req_headers.iter().map(|header| (header.0, String::from_utf8_lossy(header.1))).collect::>()); // These are the configured headers for this input. - let input_headers = input.map(|i| i.headers.clone()); // The stream url, we need to clone it because of move to async block. // We merge configured input headers with the headers from the request. let headers = get_request_headers(input_headers.as_ref(), Some(&req_headers)); - let client_builder = http_client.get(stream_url.clone()).headers(headers.clone()); + let client_builder = app_state.http_client.get(stream_url.clone()).headers(headers.clone()); let client = if stream_options.stream_connect_timeout_secs > 0 { client_builder.timeout(Duration::from_secs(u64::from(stream_options.stream_connect_timeout_secs))) } else { @@ -107,8 +105,8 @@ pub async fn get_provider_pipe_stream(cfg: &Config, let status = response.status(); if status.is_success() { (Some(Box::pin(response.bytes_stream().map_err(|err| StreamError::reqwest(&err)))), Some((response_headers, status))) - } else if let Some((boxed_provider_stream, response_info)) = create_channel_unavailable_stream(cfg, &response_headers, status) { - (Some(boxed_provider_stream), Some(response_info)) + } else if let (Some(boxed_provider_stream), response_info) = create_channel_unavailable_stream(&app_state.config, &response_headers, status) { + (Some(boxed_provider_stream), response_info) } else { (None, Some((response_headers, status))) } @@ -116,8 +114,8 @@ pub async fn get_provider_pipe_stream(cfg: &Config, Err(err) => { let masked_url = sanitize_sensitive_info(stream_url.as_str()); error!("Failed to open stream {masked_url} {err}"); - if let Some((boxed_provider_stream, response_info)) = create_channel_unavailable_stream(cfg, &get_response_headers(&headers), StatusCode::BAD_GATEWAY) { - (Some(boxed_provider_stream), Some(response_info)) + if let (Some(boxed_provider_stream), response_info) = create_channel_unavailable_stream(&app_state.config, &get_response_headers(&headers), StatusCode::BAD_GATEWAY) { + (Some(boxed_provider_stream), response_info) } else { (None, None) } @@ -125,13 +123,12 @@ pub async fn get_provider_pipe_stream(cfg: &Config, } } -pub async fn get_provider_reconnect_buffered_stream(cfg: &Config, - http_client: &Arc, +pub async fn get_provider_reconnect_buffered_stream(app_state: &AppState, stream_url: &Url, req_headers: &HeaderMap, - input: Option<&ConfigInput>, + input_headers: Option>, options: BufferStreamOptions) -> ProviderStreamResponse { - match create_provider_stream(cfg, Arc::clone(http_client), stream_url, req_headers, input, options).await { + match create_provider_stream(&app_state.config, Arc::clone(&app_state.http_client), stream_url, req_headers, input_headers, options).await { None => (None, None), Some((stream, info)) => { (Some(stream), info) diff --git a/src/api/model/streams/provider_stream_factory.rs b/src/api/model/streams/provider_stream_factory.rs index f536b5448..0a0c244c5 100644 --- a/src/api/model/streams/provider_stream_factory.rs +++ b/src/api/model/streams/provider_stream_factory.rs @@ -5,7 +5,7 @@ use crate::api::model::streams::buffered_stream::BufferedStream; use crate::api::model::streams::client_stream::ClientStream; use crate::api::model::streams::provider_stream::{create_channel_unavailable_stream, get_header_filter_for_item_type}; use crate::api::model::streams::timed_client_stream::{TimeoutClientStream}; -use crate::model::config::{Config, ConfigInput}; +use crate::model::config::{Config}; use crate::model::playlist::PlaylistItemType; use crate::tools::atomic_once_flag::AtomicOnceFlag; use crate::utils::debug_if_enabled; @@ -173,7 +173,7 @@ fn get_request_range_start_bytes(req_headers: &HashMap>) -> Opti fn get_client_stream_request_params( req_headers: &HeaderMap, - input: Option<&ConfigInput>, + input_headers: Option>, options: &BufferStreamOptions) -> (usize, Option, bool, u32, u32, HeaderMap) { let stream_buffer_size = if options.is_buffer_enabled() { options.get_stream_buffer_size() } else { 1 }; @@ -184,8 +184,6 @@ fn get_client_stream_request_params( 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. let headers = get_request_headers(input_headers.as_ref(), Some(&req_headers)); @@ -232,22 +230,22 @@ async fn provider_initial_request(cfg: &Config, request_client: Arc { - if let Some((boxed_provider_stream, response_info)) = + if let (Some(boxed_provider_stream), response_info) = create_channel_unavailable_stream(cfg, &get_response_headers(stream_options.get_headers()), StatusCode::BAD_GATEWAY) { - Ok(Some((boxed_provider_stream, Some(response_info)))) + Ok(Some((boxed_provider_stream, response_info))) } else { Err(StatusCode::SERVICE_UNAVAILABLE) } @@ -335,10 +333,10 @@ async fn get_initial_stream(cfg: &Config, client: Arc, stream_o fn create_provider_stream_options(stream_url: &Url, req_headers: &HeaderMap, - input: Option<&ConfigInput>, + input_headers: Option>, options: &BufferStreamOptions) -> ProviderStreamOptions { let (buffer_size, req_range_start_bytes, reconnect, reconnect_force_secs, connect_timeout_secs, headers) - = get_client_stream_request_params(req_headers, input, options); + = get_client_stream_request_params(req_headers, input_headers, options); let url = stream_url.clone(); let range_bytes = Arc::new(req_range_start_bytes.map(AtomicUsize::new)); let continue_flag = Arc::new(AtomicOnceFlag::new()); @@ -359,9 +357,9 @@ pub async fn create_provider_stream(cfg: &Config, client: Arc, stream_url: &Url, req_headers: &HeaderMap, - input: Option<&ConfigInput>, + input_headers: Option>, options: BufferStreamOptions) -> Option { - let stream_options = create_provider_stream_options(stream_url, req_headers, input, &options); + let stream_options = create_provider_stream_options(stream_url, req_headers, input_headers, &options); let client_stream_factory = |stream, reconnect_flag, range_cnt| { let stream = if stream_options.is_buffered() && !options.is_shared_stream() { diff --git a/src/model/api_proxy.rs b/src/model/api_proxy.rs index a7f441da8..ec826c4d3 100644 --- a/src/model/api_proxy.rs +++ b/src/model/api_proxy.rs @@ -206,6 +206,7 @@ impl ProxyUserCredentials { pub async fn has_connections_left(&self, app_state: &AppState) -> bool { if app_state.config.user_access_control { if let Some(max_connections) = self.max_connections.as_ref() { + // info!("{max_connections} : {}", app_state.get_active_connections_for_user(&self.username).await); if *max_connections < app_state.get_active_connections_for_user(&self.username).await { debug!("User access denied, too many connections: {}", self.username); return false; diff --git a/src/model/config.rs b/src/model/config.rs index 0c947de59..ebce3e029 100644 --- a/src/model/config.rs +++ b/src/model/config.rs @@ -649,6 +649,41 @@ pub struct InputUserInfo { pub password: String, } +impl InputUserInfo { + pub fn new(input_type: InputType, username: Option<&str>, password: Option<&str>, input_url: &str) -> Option { + if input_type == InputType::Xtream { + if let (Some(username), Some(password)) = (username, password) { + return Some(Self { + base_url: input_url.to_string(), + username: username.to_owned(), + password: password.to_owned(), + }); + } + } else if let Ok(url) = Url::parse(&input_url) { + let base_url = url.origin().ascii_serialization(); + let mut username = None; + let mut password = None; + for (key, value) in url.query_pairs() { + if key.eq("username") { + username = Some(value.into_owned()); + } else if key.eq("password") { + password = Some(value.into_owned()); + } + } + if username.is_some() || password.is_some() { + if let (Some(username), Some(password)) = (username.as_ref(), password.as_ref()) { + return Some(Self { + base_url, + username: username.to_owned(), + password: password.to_owned(), + }); + } + } + } + None + } +} + macro_rules! check_input_credentials { ($this:ident, $input_type:expr) => { match $input_type { @@ -756,36 +791,7 @@ impl ConfigInput { } pub fn get_user_info(&self) -> Option { - if self.input_type == InputType::Xtream { - if let (Some(username), Some(password)) = (self.username.as_ref(), self.password.as_ref()) { - return Some(InputUserInfo { - base_url: self.url.clone(), - username: username.to_owned(), - password: password.to_owned(), - }); - } - } else if let Ok(url) = Url::parse(&self.url) { - let base_url = url.origin().ascii_serialization(); - let mut username = None; - let mut password = None; - for (key, value) in url.query_pairs() { - if key.eq("username") { - username = Some(value.into_owned()); - } else if key.eq("password") { - password = Some(value.into_owned()); - } - } - if username.is_some() || password.is_some() { - if let (Some(username), Some(password)) = (username.as_ref(), password.as_ref()) { - return Some(InputUserInfo { - base_url, - username: username.to_owned(), - password: password.to_owned(), - }); - } - } - } - None + InputUserInfo::new(self.input_type.clone(), self.username.as_deref(), self.password.as_deref(), &self.url) } } @@ -1336,6 +1342,17 @@ impl Config { None } + pub fn get_input_options_by_name(&self, input_name: &str) -> Option<&ConfigInputOptions> { + for source in &self.sources { + for input in &source.inputs { + if input.name == input_name { + return input.options.as_ref(); + } + } + } + None + } + pub fn get_input_by_id(&self, input_id: u16) -> Option<&ConfigInput> { for source in &self.sources { for input in &source.inputs { diff --git a/src/model/hls.rs b/src/model/hls.rs index b9283cc85..b020f8cd9 100644 --- a/src/model/hls.rs +++ b/src/model/hls.rs @@ -5,9 +5,9 @@ use crate::model::config::TargetType; #[derive(Clone)] pub struct HlsEntry { pub ts: Instant, - pub token: String, + pub token: u32, pub target_type: TargetType, - pub input_name: String, + pub input_id: u16, pub virtual_id: u32, pub chunk: u32, pub chunks: HashMap, diff --git a/src/model/playlist.rs b/src/model/playlist.rs index 4ac1fd00d..64e4d6429 100644 --- a/src/model/playlist.rs +++ b/src/model/playlist.rs @@ -74,7 +74,7 @@ impl TryFrom for XtreamCluster { type Error = String; fn try_from(item_type: PlaylistItemType) -> Result { match item_type { - PlaylistItemType::Live | PlaylistItemType::LiveHls | PlaylistItemType::LiveUnknown => Ok(Self::Live), + PlaylistItemType::Live | PlaylistItemType::LiveHls | PlaylistItemType::LiveDash | PlaylistItemType::LiveUnknown => Ok(Self::Live), PlaylistItemType::Video => Ok(Self::Video), PlaylistItemType::Series => Ok(Self::Series), _ => Err(format!("Cant convert {item_type}")), @@ -93,6 +93,7 @@ pub enum PlaylistItemType { Catchup = 5, LiveUnknown = 6, // No Provider id LiveHls = 7, // m3u8 entry + LiveDash = 8, // mpd } impl From for PlaylistItemType { @@ -116,7 +117,7 @@ impl PlaylistItemType { impl Display for PlaylistItemType { fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result { write!(f, "{}", match self { - Self::Live | Self::LiveHls | Self::LiveUnknown => Self::LIVE, + Self::Live | Self::LiveHls | Self::LiveDash |Self::LiveUnknown => Self::LIVE, Self::Video => Self::VIDEO, Self::Series => Self::SERIES, Self::SeriesInfo => Self::SERIES_INFO, diff --git a/src/processing/parser/hls.rs b/src/processing/parser/hls.rs index b47e672b8..8bacf461d 100644 --- a/src/processing/parser/hls.rs +++ b/src/processing/parser/hls.rs @@ -4,14 +4,14 @@ use crate::model::hls::HlsEntry; use std::collections::HashMap; use std::str; use tokio::time::Instant; -use crate::utils::hash_utils::uuid_v4; use crate::utils::string_utils::replace_after_last_slash; pub const HLS_PREFIX: &str = "hls"; -pub fn rewrite_hls(base_url: &str, content: &str, hls_url: &str, virtual_id: u32, user: &ProxyUserCredentials, - target_type: &TargetType, input_name: &str) -> (HlsEntry, String) { - let token = uuid_v4(); +pub fn rewrite_hls(base_url: &str, content: &str, hls_url: &str, virtual_id: u32, + token: u32, + user: &ProxyUserCredentials, + target_type: &TargetType, input_id: u16) -> (HlsEntry, String) { let username = &user.username; let password = &user.password; let mut chunk: u32 = 1; @@ -36,7 +36,7 @@ pub fn rewrite_hls(base_url: &str, content: &str, hls_url: &str, virtual_id: u32 ts: Instant::now(), token, target_type: target_type.clone(), - input_name: input_name.to_string(), + input_id, virtual_id, chunk, chunks, diff --git a/src/processing/processor/xtream.rs b/src/processing/processor/xtream.rs index 8fd2863cc..7b2f6215d 100644 --- a/src/processing/processor/xtream.rs +++ b/src/processing/processor/xtream.rs @@ -49,7 +49,7 @@ pub(in crate::processing) fn write_info_content_to_wal_file(writer: &mut BufWrit } pub(in crate::processing) fn create_resolve_episode_wal_files(cfg: &Config, input: &ConfigInput) -> Option<(File, PathBuf)> { - match get_input_storage_path(input, &cfg.working_dir) { + match get_input_storage_path(&input.name, &cfg.working_dir) { Ok(storage_path) => { let info_path = storage_path.join(format!("{FILE_SERIES_EPISODE_RECORD}.{FILE_SUFFIX_WAL}")); let info_file = append_or_crate_file(&info_path).ok()?; @@ -60,7 +60,7 @@ pub(in crate::processing) fn create_resolve_episode_wal_files(cfg: &Config, inpu } pub(in crate::processing) fn create_resolve_info_wal_files(cfg: &Config, input: &ConfigInput, cluster: XtreamCluster) -> Option<(File, File, PathBuf, PathBuf)> { - match get_input_storage_path(input, &cfg.working_dir) { + match get_input_storage_path(&input.name, &cfg.working_dir) { Ok(storage_path) => { if let Some(file_prefix) = match cluster { XtreamCluster::Live => None, @@ -96,12 +96,12 @@ where { let mut processed_info_ids = HashMap::new(); - let file_path = match get_input_storage_path(fpl.input, &cfg.working_dir) + let fpl_name = &fpl.input.name; + let file_path = match get_input_storage_path(fpl_name, &cfg.working_dir) .map(|storage_path| xtream_get_record_file_path(&storage_path, item_type)).and_then(|opt| opt.ok_or_else(|| str_to_io_error("Not supported"))) { Ok(file_path) => file_path, Err(err) => { - let fpl_name = &fpl.input.name; errors.push(notify_err!(format!("Could not create storage path for input {fpl_name}: {err}"))); return processed_info_ids; } diff --git a/src/processing/processor/xtream_series.rs b/src/processing/processor/xtream_series.rs index 10bb47f77..0154d05db 100644 --- a/src/processing/processor/xtream_series.rs +++ b/src/processing/processor/xtream_series.rs @@ -139,7 +139,7 @@ async fn process_series_info( let mut result: Vec = vec![]; let input = fpl.input; - let Ok(Some((info_path, idx_path))) = get_input_storage_path(input, &cfg.working_dir) + let Ok(Some((info_path, idx_path))) = get_input_storage_path(&input.name, &cfg.working_dir) .map(|storage_path| xtream_get_info_file_paths(&storage_path, XtreamCluster::Series)) else { errors.push(notify_err!("Failed to open input info file for series".to_string())); diff --git a/src/repository/kodi_repository.rs b/src/repository/kodi_repository.rs index a7e7f4a8c..1f0ec438b 100644 --- a/src/repository/kodi_repository.rs +++ b/src/repository/kodi_repository.rs @@ -305,46 +305,44 @@ async fn get_tmdb_value( } } std::collections::hash_map::Entry::Vacant(entry) => { - if let Some(input) = cfg.get_input_by_name(input_name) { - 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(Some(tmdb_path)) = get_input_storage_path(input_name, &cfg.working_dir) + .map(|storage_path| xtream_get_record_file_path(&storage_path, item_type)) + { { - { - let file_lock = cfg.file_locks.read_lock(&tmdb_path).await; - match item_type { - PlaylistItemType::Series => { - if let Ok(tree) = - BPlusTree::::load(&tmdb_path) - { - let tmdb_id = tree.query(&pid).map(|episode| { - InputTmdbIndexValue::Series(episode.clone()) - }); - entry.insert(Some(( - file_lock, - InputTmdbIndexTree::Series(tree), - ))); - return tmdb_id; - } + let file_lock = cfg.file_locks.read_lock(&tmdb_path).await; + match item_type { + PlaylistItemType::Series => { + if let Ok(tree) = + BPlusTree::::load(&tmdb_path) + { + let tmdb_id = tree.query(&pid).map(|episode| { + InputTmdbIndexValue::Series(episode.clone()) + }); + entry.insert(Some(( + file_lock, + InputTmdbIndexTree::Series(tree), + ))); + return tmdb_id; } - PlaylistItemType::Video => { - if let Ok(tree) = - BPlusTree::::load(&tmdb_path) - { - let tmdb_id = tree.query(&pid).map(|vod_record| { - InputTmdbIndexValue::Video(vod_record.clone()) - }); - entry.insert(Some(( - file_lock, - InputTmdbIndexTree::Video(tree), - ))); - return tmdb_id; - } - } - _ => {} } - }; + PlaylistItemType::Video => { + if let Ok(tree) = + BPlusTree::::load(&tmdb_path) + { + let tmdb_id = tree.query(&pid).map(|vod_record| { + InputTmdbIndexValue::Video(vod_record.clone()) + }); + entry.insert(Some(( + file_lock, + InputTmdbIndexTree::Video(tree), + ))); + return tmdb_id; + } + } + _ => {} + } }; - } + }; entry.insert(None); None } diff --git a/src/repository/m3u_playlist_iterator.rs b/src/repository/m3u_playlist_iterator.rs index 175e34ad2..ee204f6b5 100644 --- a/src/repository/m3u_playlist_iterator.rs +++ b/src/repository/m3u_playlist_iterator.rs @@ -70,6 +70,7 @@ impl M3uPlaylistIterator { | PlaylistItemType::Catchup | PlaylistItemType::LiveUnknown | PlaylistItemType::LiveHls => "live", + | PlaylistItemType::LiveDash => "live", PlaylistItemType::Video => "movie", PlaylistItemType::Series | PlaylistItemType::SeriesInfo => "series", diff --git a/src/repository/playlist_repository.rs b/src/repository/playlist_repository.rs index 900e7d23b..cc0fa6900 100644 --- a/src/repository/playlist_repository.rs +++ b/src/repository/playlist_repository.rs @@ -1,5 +1,4 @@ -use std::path::{Path}; -use crate::m3u_filter_error::{info_err}; +use crate::m3u_filter_error::info_err; use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; use crate::model::config::{Config, ConfigTarget, TargetOutput}; use crate::model::playlist::{PlaylistGroup, PlaylistItemType}; @@ -11,7 +10,8 @@ use crate::repository::storage::{ensure_target_storage_path, get_target_id_mappi use crate::repository::target_id_mapping::TargetIdMapping; use crate::repository::xtream_repository::xtream_write_playlist; use crate::utils::file::file_lock_manager::FileWriteGuard; -use crate::utils::network::request::is_hls_url; +use crate::utils::network::request::{is_dash_url, is_hls_url}; +use std::path::Path; pub async fn persist_playlist(playlist: &mut [PlaylistGroup], epg: Option<&Epg>, target: &ConfigTarget, cfg: &Config) -> Result<(), Vec> { @@ -31,7 +31,13 @@ pub async fn persist_playlist(playlist: &mut [PlaylistGroup], epg: Option<&Epg>, if provider_id == 0 { header.item_type = match (is_hls_url(&header.url), header.item_type) { (true, _) => PlaylistItemType::LiveHls, - (false, PlaylistItemType::Live) => PlaylistItemType::LiveUnknown, + (false, PlaylistItemType::Live) => { + if is_dash_url(&header.url) { + PlaylistItemType::LiveDash + } else { + PlaylistItemType::LiveUnknown + } + } _ => header.item_type, }; } diff --git a/src/repository/storage.rs b/src/repository/storage.rs index 8d1c273e7..8ac586f61 100644 --- a/src/repository/storage.rs +++ b/src/repository/storage.rs @@ -1,7 +1,7 @@ use std::path::{Path, PathBuf}; use std::fmt::Write; use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; -use crate::model::config::{Config, ConfigInput}; +use crate::model::config::{Config}; use crate::model::playlist::UUIDType; use crate::m3u_filter_error::{notify_err}; use crate::utils::file::file_utils; @@ -56,8 +56,8 @@ pub fn get_target_storage_path(cfg: &Config, target_name: &str) -> Option std::io::Result { - let name = format!("input_{}", &input.name); +pub fn get_input_storage_path(input_name: &str, working_dir: &str) -> std::io::Result { + let name = format!("input_{input_name}"); let path = Path::new(working_dir).join(name); // Create the directory and return the path or propagate the error std::fs::create_dir_all(&path).map(|()| path) diff --git a/src/repository/xtream_repository.rs b/src/repository/xtream_repository.rs index 1ef47152e..b193c2314 100644 --- a/src/repository/xtream_repository.rs +++ b/src/repository/xtream_repository.rs @@ -744,7 +744,7 @@ pub fn xtream_get_input_info( provider_id: u32, cluster: XtreamCluster, ) -> 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(Some((info_path, idx_path))) = get_input_storage_path(&input.name, &cfg.working_dir).map(|storage_path| xtream_get_info_file_paths(&storage_path, cluster)) { let _file_lock = cfg.file_locks.read_lock(&info_path); if let Ok(content) = IndexedDocumentDirectAccess::read_indexed_item::(&info_path, &idx_path, &provider_id) { @@ -760,7 +760,7 @@ pub async fn xtream_update_input_info_file( wal_path: &Path, cluster: XtreamCluster, ) -> Result<(), M3uFilterError> { - match get_input_storage_path(input, &cfg.working_dir).map(|storage_path| xtream_get_info_file_paths(&storage_path, cluster)) { + match get_input_storage_path(&input.name, &cfg.working_dir).map(|storage_path| xtream_get_info_file_paths(&storage_path, cluster)) { Ok(Some((info_path, idx_path))) => { { let _file_lock = cfg.file_locks.write_lock(&info_path); @@ -803,7 +803,7 @@ pub async fn xtream_update_input_vod_record_from_wal_file( input: &ConfigInput, wal_path: &Path, ) -> Result<(), M3uFilterError> { - let record_path = get_input_storage_path(input, &cfg.working_dir).map(|storage_path| xtream_get_record_file_path(&storage_path, PlaylistItemType::Video)) + let record_path = get_input_storage_path(&input.name, &cfg.working_dir).map(|storage_path| xtream_get_record_file_path(&storage_path, PlaylistItemType::Video)) .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))))?; @@ -843,7 +843,7 @@ pub async fn xtream_update_input_series_record_from_wal_file( input: &ConfigInput, wal_path: &Path, ) -> Result<(), M3uFilterError> { - let record_path = get_input_storage_path(input, &cfg.working_dir).map(|storage_path| xtream_get_record_file_path(&storage_path, PlaylistItemType::SeriesInfo)) + let record_path = get_input_storage_path(&input.name, &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))))?; { @@ -877,7 +877,7 @@ pub async fn xtream_update_input_series_episodes_record_from_wal_file( input: &ConfigInput, wal_path: &Path, ) -> Result<(), M3uFilterError> { - let record_path = get_input_storage_path(input, &cfg.working_dir).map(|storage_path| xtream_get_record_file_path(&storage_path, PlaylistItemType::Series)) + let record_path = get_input_storage_path(&input.name, &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))))?; { diff --git a/src/utils/hash_utils.rs b/src/utils/hash_utils.rs index d5d2026c1..7adf72410 100644 --- a/src/utils/hash_utils.rs +++ b/src/utils/hash_utils.rs @@ -22,8 +22,3 @@ pub fn generate_playlist_uuid(key: &str, provider_id: &str, item_type: PlaylistI } hash_string(url) } - -#[inline] -pub fn uuid_v4() -> String { - uuid::Uuid::new_v4().to_string().replace('-', "") -} diff --git a/src/utils/network/request.rs b/src/utils/network/request.rs index 7221dd5ef..a4d612dd8 100644 --- a/src/utils/network/request.rs +++ b/src/utils/network/request.rs @@ -27,6 +27,7 @@ use crate::m3u_filter_error::create_m3u_filter_error_result; use crate::utils::debug_if_enabled; pub const HLS_EXT: &str = ".m3u8"; +pub const DASH_EXT: &str = ".mpd"; pub const fn bytes_to_megabytes(bytes: u64) -> u64 { bytes / 1_048_576 @@ -289,7 +290,7 @@ pub async fn download_text_content_as_file(client: Arc, input: Err(Error::new(ErrorKind::NotFound, format!("Unknown file {file_path:?}"))) }) } else { - let file_path = persist_filepath.map_or_else(|| match get_input_storage_path(input, working_dir) { + let file_path = persist_filepath.map_or_else(|| match get_input_storage_path(&input.name, working_dir) { Ok(download_path) => { Ok(download_path.join(FILE_EPG)) } @@ -430,12 +431,20 @@ pub fn classify_content_type(headers: &[(String, String)]) -> MimeCategory { const HLS_EXT_QUERY: &str = ".m3u8?"; const HLS_EXT_FRAGMENT: &str = ".m3u8#"; +const DASH_EXT_QUERY: &str = ".mpd?"; +const DASH_EXT_FRAGMENT: &str = ".mpd#"; + pub fn is_hls_url(url: &str) -> bool { let lc_url = url.to_lowercase(); lc_url.ends_with(HLS_EXT) || lc_url.contains(HLS_EXT_QUERY) || lc_url.contains(HLS_EXT_FRAGMENT) } +pub fn is_dash_url(url: &str) -> bool { + let lc_url = url.to_lowercase(); + lc_url.ends_with(DASH_EXT) || lc_url.contains(DASH_EXT_QUERY) || lc_url.contains(DASH_EXT_FRAGMENT) +} + pub fn replace_url_extension(url: &str, new_ext: &str) -> String { let ext = new_ext.strip_prefix('.').unwrap_or(new_ext); // Remove leading dot if exists