From 9d35a9cbbe58fcdc751c99bef3e3a0048ce0d6ea Mon Sep 17 00:00:00 2001 From: euzu <33094714+euzu@users.noreply.github.com> Date: Fri, 27 Mar 2026 15:23:39 +0100 Subject: [PATCH] Sorce ordinal by category and channel index (#672) * Sorce ordinal by category and channel index * Parsing UserApiRequest from request body or query * CSS stylings * media tools refactoring * xtream info parsing now number/string fields --- Cargo.lock | 117 +++-- Cargo.toml | 6 +- backend/Cargo.toml | 11 +- backend/src/api/api_utils.rs | 2 +- backend/src/api/endpoints/m3u_api.rs | 7 +- backend/src/api/endpoints/xmltv_api.rs | 10 +- backend/src/api/endpoints/xtream_api.rs | 17 +- .../src/api/model/active_provider_manager.rs | 4 +- backend/src/api/model/app_state.rs | 2 +- backend/src/api/model/request.rs | 293 +++++++++++- .../api/model/streams/active_client_stream.rs | 4 +- .../src/api/model/streams/provider_stream.rs | 4 +- .../model/streams/shared_stream_manager.rs | 4 +- backend/src/library/processor.rs | 73 ++- backend/src/library/thumbnail.rs | 97 +--- backend/src/messaging.rs | 4 +- backend/src/model/config/app.rs | 14 +- backend/src/model/config/media_tools.rs | 48 ++ backend/src/model/config/mod.rs | 2 + backend/src/modules.rs | 1 - backend/src/processing/parser/xtream.rs | 236 +++++++++- .../processor/probe_handle_guard.rs | 4 +- .../src/processing/processor/stream_probe.rs | 6 +- backend/src/processing/processor/xtream.rs | 4 +- .../src/processing/processor/xtream_series.rs | 4 +- .../src/processing/processor/xtream_vod.rs | 4 +- backend/src/repository/user_repository.rs | 3 +- backend/src/tools/mod.rs | 1 - backend/src/utils/ffmpeg.rs | 441 ++++++++++++------ backend/src/utils/file/config_reader.rs | 5 +- backend/src/{tools => utils}/lru_cache.rs | 0 backend/src/utils/mod.rs | 2 + backend/src/utils/network/request.rs | 4 +- docs/src/configuration/api-proxy.md | 321 ++++++++++++- frontend/public/assets/i18n/en.json | 3 +- frontend/scss/_theme.scss | 27 +- .../app/components/_radio_button_group.scss | 18 +- frontend/scss/app/components/_tabset.scss | 8 +- .../scss/app/components/_text_button.scss | 24 +- .../app/components/dashboard/_stats_view.scss | 14 + frontend/src/app/components/card.rs | 4 +- .../app/components/dashboard/stats_view.rs | 22 +- .../playlist/playlist_update_view.rs | 16 +- frontend/src/services/config_service.rs | 4 +- shared/src/model/xtream.rs | 27 ++ shared/src/utils/string_interner.rs | 143 ++++-- 46 files changed, 1644 insertions(+), 421 deletions(-) create mode 100644 backend/src/model/config/media_tools.rs delete mode 100644 backend/src/tools/mod.rs rename backend/src/{tools => utils}/lru_cache.rs (100%) diff --git a/Cargo.lock b/Cargo.lock index 94a255947..d788680f6 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -73,9 +73,9 @@ dependencies = [ [[package]] name = "anstream" -version = "0.6.21" +version = "1.0.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "43d5b281e737544384e969a5ccad3f1cdd24b48086a0fc1b2a5262a26b8f4f4a" +checksum = "824a212faf96e9acacdbd09febd34438f8f711fb84e09a8916013cd7815ca28d" dependencies = [ "anstyle", "anstyle-parse", @@ -94,9 +94,9 @@ checksum = "5192cca8006f1fd4f7237516f40fa183bb07f8fbdfedaa0036de5ea9b0b45e78" [[package]] name = "anstyle-parse" -version = "0.2.7" +version = "1.0.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "4e7644824f0aa2c7b9384579234ef10eb7efb6a0deb83f9630a49594dd9c15c2" +checksum = "52ce7f38b242319f7cabaa6813055467063ecdc9d355bbb4ce0c68908cd8130e" dependencies = [ "utf8parse", ] @@ -129,9 +129,9 @@ checksum = "7f202df86484c868dbad7eaa557ef785d5c66295e41b460ef922eca0723b842c" [[package]] name = "arc-swap" -version = "1.8.2" +version = "1.9.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "f9f3647c145568cec02c42054e07bdf9a5a698e15b466fb2341bfc393cd24aa5" +checksum = "a07d1f37ff60921c83bdfc7407723bdefe89b44b98a9b772f225c8f9d67141a6" dependencies = [ "rustversion", ] @@ -526,9 +526,9 @@ dependencies = [ [[package]] name = "clap" -version = "4.5.60" +version = "4.6.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "2797f34da339ce31042b27d23607e051786132987f595b02ba4f6a6dffb7030a" +checksum = "b193af5b67834b676abd72466a96c1024e6a6ad978a1f484bd90b85c94041351" dependencies = [ "clap_builder", "clap_derive", @@ -536,9 +536,9 @@ dependencies = [ [[package]] name = "clap_builder" -version = "4.5.60" +version = "4.6.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "24a241312cea5059b13574bb9b3861cabf758b879c15190b37b6d6fd63ab6876" +checksum = "714a53001bf66416adb0e2ef5ac857140e7dc3a0c48fb28b2f10762fc4b5069f" dependencies = [ "anstream", "anstyle", @@ -548,9 +548,9 @@ dependencies = [ [[package]] name = "clap_derive" -version = "4.5.55" +version = "4.6.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "a92793da1a46a5f2a02a6f4c46c6496b28c43638adea8306fcb0caa1634f24e5" +checksum = "1110bd8a634a1ab8cb04345d8d878267d57c3cf1b38d91b71af6686408bbca6a" dependencies = [ "heck", "proc-macro2", @@ -677,13 +677,14 @@ dependencies = [ [[package]] name = "cron" -version = "0.15.0" +version = "0.16.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "5877d3fbf742507b66bc2a1945106bd30dd8504019d596901ddd012a4dd01740" +checksum = "089df96cf6a25253b4b6b6744d86f91150a3d4df546f31a95def47976b8cba97" dependencies = [ "chrono", "once_cell", - "winnow 0.6.26", + "phf 0.11.3", + "winnow 0.7.15", ] [[package]] @@ -1110,9 +1111,9 @@ dependencies = [ [[package]] name = "env_logger" -version = "0.11.9" +version = "0.11.10" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "b2daee4ea451f429a58296525ddf28b45a3b64f1acf6587e2067437bb11e218d" +checksum = "0621c04f2196ac3f488dd583365b9c09be011a4ab8b9f37248ffcc8f6198b56a" dependencies = [ "anstream", "anstyle", @@ -2732,9 +2733,9 @@ checksum = "d6790f58c7ff633d8771f42965289203411a5e5c68388703c06e14f24770b41e" [[package]] name = "openssl" -version = "0.10.75" +version = "0.10.76" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "08838db121398ad17ab8531ce9de97b244589089e290a384c900cb9ff7434328" +checksum = "951c002c75e16ea2c65b8c7e4d3d51d5530d8dfa7d060b4776828c88cfb18ecf" dependencies = [ "bitflags 2.11.0", "cfg-if", @@ -2773,9 +2774,9 @@ dependencies = [ [[package]] name = "openssl-sys" -version = "0.9.111" +version = "0.9.112" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "82cab2d520aa75e3c58898289429321eb788c3106963d0dc886ec7a5f4adc321" +checksum = "57d55af3b3e226502be1526dfdba67ab0e9c96fc293004e79576b2b9edb0dbdb" dependencies = [ "cc", "libc", @@ -2911,6 +2912,16 @@ dependencies = [ "sha2", ] +[[package]] +name = "phf" +version = "0.11.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1fd6780a80ae0c52cc120a26a1a42c1ae51b247a253e4e06113d23d2c2edd078" +dependencies = [ + "phf_macros 0.11.3", + "phf_shared 0.11.3", +] + [[package]] name = "phf" version = "0.12.1" @@ -2926,7 +2937,7 @@ version = "0.13.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c1562dc717473dbaa4c1f85a36410e03c047b2e7df7f45ee938fbef64ae7fadf" dependencies = [ - "phf_macros", + "phf_macros 0.13.1", "phf_shared 0.13.1", "serde", ] @@ -2937,10 +2948,20 @@ version = "0.13.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "49aa7f9d80421bca176ca8dbfebe668cc7a2684708594ec9f3c0db0805d5d6e1" dependencies = [ - "phf_generator", + "phf_generator 0.13.1", "phf_shared 0.13.1", ] +[[package]] +name = "phf_generator" +version = "0.11.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3c80231409c20246a13fddb31776fb942c38553c51e871f8cbd687a4cfb5843d" +dependencies = [ + "phf_shared 0.11.3", + "rand 0.8.5", +] + [[package]] name = "phf_generator" version = "0.13.1" @@ -2951,19 +2972,41 @@ dependencies = [ "phf_shared 0.13.1", ] +[[package]] +name = "phf_macros" +version = "0.11.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f84ac04429c13a7ff43785d75ad27569f2951ce0ffd30a3321230db2fc727216" +dependencies = [ + "phf_generator 0.11.3", + "phf_shared 0.11.3", + "proc-macro2", + "quote", + "syn 2.0.117", +] + [[package]] name = "phf_macros" version = "0.13.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "812f032b54b1e759ccd5f8b6677695d5268c588701effba24601f6932f8269ef" dependencies = [ - "phf_generator", + "phf_generator 0.13.1", "phf_shared 0.13.1", "proc-macro2", "quote", "syn 2.0.117", ] +[[package]] +name = "phf_shared" +version = "0.11.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "67eabc2ef2a60eb7faa00097bd1ffdb5bd28e62bf39990626a582201b7a754e5" +dependencies = [ + "siphasher", +] + [[package]] name = "phf_shared" version = "0.12.1" @@ -3723,9 +3766,9 @@ dependencies = [ [[package]] name = "saphyr-parser-bw" -version = "0.0.608" +version = "0.0.610" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "d55ae5ea09894b6d5382621db78f586df37ef18ab581bf32c754e75076b124b1" +checksum = "6d643f5e972f17219245b82f038c22cd3c74320bb17c6e8f7e8537de268b1bc6" dependencies = [ "arraydeque", "smallvec", @@ -3821,9 +3864,9 @@ dependencies = [ [[package]] name = "serde-saphyr" -version = "0.0.21" +version = "0.0.22" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "4a6fc4aa0da972ba0f51cf5c1bb16e9dba35334adc6831b09b3ffb0ec20bb264" +checksum = "546b4da4f679832602a8f8ab8ddc10b6b1d2e1a13b4f9dddcaee499436fa06ad" dependencies = [ "ahash", "annotate-snippets", @@ -3870,6 +3913,19 @@ dependencies = [ "syn 2.0.117", ] +[[package]] +name = "serde_html_form" +version = "0.4.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0946d52b4b7e28823148aebbeceb901012c595ad737920d504fa8634bb099e6f" +dependencies = [ + "form_urlencoded", + "indexmap", + "itoa", + "ryu", + "serde_core", +] + [[package]] name = "serde_json" version = "1.0.149" @@ -4607,6 +4663,7 @@ dependencies = [ "rustls", "serde", "serde-saphyr", + "serde_html_form", "serde_json", "shared", "socket2", @@ -5400,9 +5457,9 @@ dependencies = [ [[package]] name = "winnow" -version = "0.6.26" +version = "0.7.15" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1e90edd2ac1aa278a5c4599b1d89cf03074b610800f866d4026dc199d7929a28" +checksum = "df79d97927682d2fd8adb29682d1140b343be4ac0f08fd68b7765d9c059d3945" dependencies = [ "memchr", ] diff --git a/Cargo.toml b/Cargo.toml index 0ca333d64..3989b6b23 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -20,7 +20,7 @@ fastrand = "2.3.0" lz4_flex = "0.13.0" indexmap = "2.13.0" dashmap = "6.1.0" -arc-swap = "1.8.2" +arc-swap = "1.9.0" paste = "1.0.15" pest = "2.8.6" pest_derive = "2.8.6" @@ -28,10 +28,10 @@ zeroize = "1.8.2" deunicode = "1.6.2" path-clean = "1.0.1" url = "2.5.8" -cron = "0.15.0" +cron = "0.16.0" enum-iterator = "2.3.0" futures = "0.3.32" -serde-saphyr = "0.0.21" +serde-saphyr = "0.0.22" tokio = { version = "1.50.0" } [profile.release] diff --git a/backend/Cargo.toml b/backend/Cargo.toml index 8530aec84..b20d4f89c 100644 --- a/backend/Cargo.toml +++ b/backend/Cargo.toml @@ -15,13 +15,14 @@ serde-saphyr.workspace = true serde_json = { workspace = true, features = ["raw_value"] } quick-xml = { version = "0.39.2", features = ["async-tokio", "serialize"] } regex.workspace = true -clap = { version = "4.5.60", features = ["derive"] } +clap = { version = "4.6.0", features = ["derive"] } reqwest = { version = "0.13.2", features = ["json", "stream", "rustls", "socks", "form"] } url.workspace = true chrono.workspace = true iana-time-zone = "0.1.65" cron.workspace = true axum = { version = "0.8.8" , features = ["macros", "default", "ws"]} +serde_html_form = "0.4.0" tower = "0.5.3" tower-http = { version = "0.6.8", features = ["cors", "auth", "fs", "compression-full", "trace"] } tower_governor = { version = "0.8.0", features = ["axum"] } @@ -32,11 +33,11 @@ path-clean.workspace = true pest.workspace = true pest_derive.workspace = true enum-iterator.workspace = true -openssl = { version = "0.10.75", features = ["vendored"] } #https://docs.rs/openssl/0.10.34/openssl/#vendored +openssl = { version = "0.10.76", features = ["vendored"] } #https://docs.rs/openssl/0.10.34/openssl/#vendored deunicode.workspace = true mime = "0.3.17" log.workspace = true -env_logger = "0.11.9" +env_logger = "0.11.10" rmp-serde = "1.3.1" rand = "0.9.2" fastrand.workspace = true @@ -75,7 +76,8 @@ uuid = { version = "1.22.0", features = ["v4"] } fancy-regex = "0.17.0" mime_guess = "2.0.5" crc32fast = "1.5.0" -memchr = "2.7.6" +memchr = "2.8.0" +http-body-util = "0.1.3" [target.'cfg(unix)'.dependencies] libc = "0.2" @@ -94,7 +96,6 @@ windows-sys = { version = "0.61.2", features = ["Win32_Foundation", "Win32_Syste vergen = { version = "9.1.0", features = ["build"] } [dev-dependencies] -http-body-util = "0.1.3" tokio = { workspace = true, features = ["test-util", "macros"] } rcgen = "0.14.7" rustls = "0.23.37" diff --git a/backend/src/api/api_utils.rs b/backend/src/api/api_utils.rs index d51d6d00d..21adb5926 100644 --- a/backend/src/api/api_utils.rs +++ b/backend/src/api/api_utils.rs @@ -11,7 +11,6 @@ use crate::{ }, auth::Fingerprint, model::{ConfigInput, ConfigTarget, ProxyUserCredentials}, - tools::lru_cache::LRUResourceCache, utils::{ async_file_reader, async_file_writer, create_new_file_for_write, debug_if_enabled, get_file_extension, request, request::{content_type_from_ext, parse_range, send_with_retry_and_provider}, @@ -206,6 +205,7 @@ pub use try_result_bad_request; pub use try_result_not_found; pub use try_result_or_status; pub use try_unwrap_body; +use crate::utils::LRUResourceCache; pub fn get_server_time() -> String { chrono::offset::Local::now().with_timezone(&chrono::Local).format("%Y-%m-%d %H:%M:%S %Z").to_string() diff --git a/backend/src/api/endpoints/m3u_api.rs b/backend/src/api/endpoints/m3u_api.rs index 0349bfe7e..b48f126f7 100644 --- a/backend/src/api/endpoints/m3u_api.rs +++ b/backend/src/api/endpoints/m3u_api.rs @@ -12,7 +12,7 @@ use crate::{ hls_api::handle_hls_stream_request, xtream_api::{ApiStreamContext, ApiStreamRequest}, }, - model::{create_custom_video_stream_response, AppState, CustomVideoStreamType, UserApiRequest}, + model::{create_custom_video_stream_response, AppState, CustomVideoStreamType, UserApiRequestQueryOrBody, UserApiRequest}, }, auth::Fingerprint, repository::{m3u_get_item_for_stream_id, m3u_load_rewrite_playlist, storage_const}, @@ -69,12 +69,9 @@ async fn m3u_api_get( } async fn m3u_api_post( - axum::extract::Query(api_query_req): axum::extract::Query, axum::extract::State(app_state): axum::extract::State>, - api_form_req: Result, axum::extract::rejection::FormRejection>, + UserApiRequestQueryOrBody(api_req): UserApiRequestQueryOrBody, ) -> impl IntoResponse + Send { - let form_req = api_form_req.as_ref().ok().map(|form| &form.0); - let api_req = UserApiRequest::merge_query_over_form(&api_query_req, form_req); m3u_api(&api_req, &app_state).await.into_response() } diff --git a/backend/src/api/endpoints/xmltv_api.rs b/backend/src/api/endpoints/xmltv_api.rs index 5b0d804ca..7eb05fa22 100644 --- a/backend/src/api/endpoints/xmltv_api.rs +++ b/backend/src/api/endpoints/xmltv_api.rs @@ -4,7 +4,7 @@ use crate::{ empty_json_response_as_array, get_user_target, get_user_target_by_credentials, internal_server_error, resource_response, stream_json_or_bin_response_stream, try_option_forbidden, try_unwrap_body, }, - model::{AppState, UserApiRequest}, + model::{AppState, UserApiRequestQueryOrBody, UserApiRequest}, }, model::{Config, ConfigTarget, ProxyUserCredentials, TargetOutput, EPG_ATTRIB_ID, EPG_TAG_CHANNEL}, repository::{ @@ -492,14 +492,8 @@ async fn xmltv_api_get( async fn xmltv_api_post( axum::extract::State(app_state): axum::extract::State>, - axum::extract::Query(api_query_req): axum::extract::Query, - api_form_req: Result, axum::extract::rejection::FormRejection>, + UserApiRequestQueryOrBody(api_req): UserApiRequestQueryOrBody, ) -> impl IntoResponse + Send { - if let Err(ref rejection) = api_form_req { - debug!("xmltv_api_post: form parsing failed: {rejection:?}"); - } - let form_req = api_form_req.as_ref().ok().map(|form| &form.0); - let api_req = UserApiRequest::merge_query_over_form(&api_query_req, form_req); xmltv_api(api_req, &app_state).await } diff --git a/backend/src/api/endpoints/xtream_api.rs b/backend/src/api/endpoints/xtream_api.rs index a991f923f..4ad0418b6 100644 --- a/backend/src/api/endpoints/xtream_api.rs +++ b/backend/src/api/endpoints/xtream_api.rs @@ -17,7 +17,7 @@ use crate::{ xmltv_api::{get_empty_epg_response, get_epg_path_for_target, serve_short_epg}, }, model::{ - create_custom_video_stream_response, AppState, CustomVideoStreamType, UserApiRequest, + create_custom_video_stream_response, AppState, CustomVideoStreamType, UserApiRequestQueryOrBody, UserApiRequest, XtreamAuthorizationResponse, }, }, @@ -698,13 +698,10 @@ struct XtreamApiTimeShiftRequest { async fn xtream_player_api_timeshift_stream( fingerprint: Fingerprint, req_headers: HeaderMap, - axum::extract::Query(api_query_req): axum::extract::Query, axum::extract::Path(timeshift_request): axum::extract::Path, axum::extract::State(app_state): axum::extract::State>, - api_form_req: Result, axum::extract::rejection::FormRejection>, + UserApiRequestQueryOrBody(query_req): UserApiRequestQueryOrBody, ) -> impl IntoResponse + Send { - let form_req = api_form_req.as_ref().ok().map(|form| &form.0); - let query_req = UserApiRequest::merge_query_over_form(&api_query_req, form_req); let path_req = UserApiRequest { username: timeshift_request.username, password: timeshift_request.password, @@ -752,12 +749,9 @@ async fn xtream_player_api_timeshift_stream( async fn xtream_player_api_timeshift_query_stream( fingerprint: Fingerprint, req_headers: HeaderMap, - axum::extract::Query(api_query_req): axum::extract::Query, axum::extract::State(app_state): axum::extract::State>, - api_form_req: Result, axum::extract::rejection::FormRejection>, + UserApiRequestQueryOrBody(api_req): UserApiRequestQueryOrBody, ) -> impl IntoResponse + Send { - let form_req = api_form_req.as_ref().ok().map(|form| &form.0); - let api_req = UserApiRequest::merge_query_over_form(&api_query_req, form_req); if api_req.username.is_empty() || api_req.password.is_empty() @@ -1314,11 +1308,8 @@ async fn xtream_player_api_get( async fn xtream_player_api_post( axum::extract::State(app_state): axum::extract::State>, - axum::extract::Query(api_query_req): axum::extract::Query, - api_form_req: Result, axum::extract::rejection::FormRejection>, + UserApiRequestQueryOrBody(api_req): UserApiRequestQueryOrBody, ) -> impl IntoResponse + Send { - let form_req = api_form_req.as_ref().ok().map(|form| &form.0); - let api_req = UserApiRequest::merge_query_over_form(&api_query_req, form_req); xtream_player_api(api_req, &app_state).await } diff --git a/backend/src/api/model/active_provider_manager.rs b/backend/src/api/model/active_provider_manager.rs index 352a82e6d..d6a596a77 100644 --- a/backend/src/api/model/active_provider_manager.rs +++ b/backend/src/api/model/active_provider_manager.rs @@ -1104,7 +1104,7 @@ mod tests { use super::ActiveProviderManager; use crate::{ api::model::{EventManager, ProviderAllocation}, - model::{AppConfig, Config, ConfigInput, ConfigInputAlias, SourcesConfig}, + model::{AppConfig, Config, ConfigInput, ConfigInputAlias, MediaToolCapabilities, SourcesConfig}, utils::FileLockManager, }; use arc_swap::{ArcSwap, ArcSwapOption}; @@ -1156,7 +1156,7 @@ mod tests { custom_stream_response: Arc::new(ArcSwapOption::default()), access_token_secret: [0; 32], encrypt_secret: [0; 16], - ffprobe_available: Arc::default(), + media_tools: Arc::new(MediaToolCapabilities::new()), } } diff --git a/backend/src/api/model/app_state.rs b/backend/src/api/model/app_state.rs index 77851946a..03f1ceae0 100644 --- a/backend/src/api/model/app_state.rs +++ b/backend/src/api/model/app_state.rs @@ -14,7 +14,6 @@ use crate::{ ProcessTargets, ReverseProxyDisabledHeaderConfig, ScheduleConfig, SourcesConfig, }, repository::{get_geoip_path, load_target_into_memory_cache}, - tools::lru_cache::LRUResourceCache, utils::{ request::{create_client, create_client_with_redirect}, GeoIp, @@ -37,6 +36,7 @@ use tokio::sync::{mpsc, Mutex}; use tokio::task; use tokio_util::sync::CancellationToken; use url::Url; +use crate::utils::LRUResourceCache; macro_rules! cancel_service { ($field: ident, $flag:expr, $changes:expr, $cancel_tokens:expr) => { diff --git a/backend/src/api/model/request.rs b/backend/src/api/model/request.rs index 0abdea015..7c3c0fdbf 100644 --- a/backend/src/api/model/request.rs +++ b/backend/src/api/model/request.rs @@ -1,4 +1,11 @@ +use axum::extract::{FromRequest, Request}; +use axum::http::StatusCode; +use axum::response::{IntoResponse, Response}; +use http_body_util::LengthLimitError; use log::log_enabled; +use std::error::Error as StdError; + +const MAX_BODY_SIZE_BYTES: usize = 10 * 1024 * 1024; #[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize, Default)] pub struct UserApiRequest { @@ -32,7 +39,194 @@ pub struct UserApiRequest { pub content_type: String, } +/// Custom extractor that parses `UserApiRequest` from query parameters and request body +/// (either `application/x-www-form-urlencoded` or `multipart/form-data`), then merges +/// them with query parameters taking priority over body fields. +pub struct UserApiRequestQueryOrBody(pub UserApiRequest); + +#[derive(Debug)] +enum ParseBodyError { + PayloadTooLarge(String), + BadRequest(String), +} + +impl ParseBodyError { + fn into_response(self) -> Response { + match self { + Self::PayloadTooLarge(err) => (StatusCode::PAYLOAD_TOO_LARGE, err).into_response(), + Self::BadRequest(err) => (StatusCode::BAD_REQUEST, err).into_response(), + } + } +} + +impl FromRequest for UserApiRequestQueryOrBody +where + S: Send + Sync, +{ + type Rejection = Response; + + async fn from_request(req: Request, _state: &S) -> Result { + let (parts, body) = req.into_parts(); + + let query_req = parse_query_request(parts.uri.query()) + .map_err(ParseBodyError::into_response)?; + + let content_type = parts + .headers + .get("content-type") + .and_then(|v| v.to_str().ok()) + .unwrap_or(""); + + let body_req = match parse_body(body, content_type).await { + Ok(body_req) => Some(body_req), + Err(err) => return Err(err.into_response()), + }; + + Ok(UserApiRequestQueryOrBody(UserApiRequest::merge_query_over_form(&query_req, body_req.as_ref()))) + } +} + +async fn parse_body(body: axum::body::Body, content_type: &str) -> Result { + let bytes = axum::body::to_bytes(body, MAX_BODY_SIZE_BYTES) + .await + .map_err(|e| { + if is_length_limit_error(&e) { + ParseBodyError::PayloadTooLarge(format!("Request body too large (max {MAX_BODY_SIZE_BYTES} bytes)")) + } else { + ParseBodyError::BadRequest(format!("Failed to read request body: {e}")) + } + })?; + + if bytes.is_empty() { + return Ok(UserApiRequest::default()); + } + + if is_multipart_content_type(content_type) { + parse_multipart_body(&bytes, content_type) + } else { + // Treat as form-urlencoded (works for both explicit content-type and missing content-type) + serde_html_form::from_bytes(&bytes).map_err(|e| ParseBodyError::BadRequest(format!("Failed to parse form: {e}"))) + } +} + +fn parse_multipart_body(bytes: &[u8], content_type: &str) -> Result { + let boundary = extract_multipart_boundary(content_type) + .ok_or_else(|| ParseBodyError::BadRequest("Missing boundary in multipart content type".to_string()))?; + + let data = std::str::from_utf8(bytes).map_err(|e| ParseBodyError::BadRequest(format!("Invalid UTF-8: {e}")))?; + + let mut request = UserApiRequest::default(); + let delimiter = format!("--{boundary}"); + for part in data.split(&delimiter) { + if let Some((name, value)) = parse_multipart_field(part) { + request.set_field(name, value); + } + } + Ok(request) +} + +fn parse_query_request(query: Option<&str>) -> Result { + match query { + Some(query) => serde_html_form::from_str(query) + .map_err(|e| ParseBodyError::BadRequest(format!("Failed to parse query: {e}"))), + None => Ok(UserApiRequest::default()), + } +} + +fn parse_multipart_field(part: &str) -> Option<(&str, &str)> { + let header_end = part.find("\r\n\r\n")?; + let headers = &part[..header_end]; + let body = &part[header_end + 4..]; + let name = extract_multipart_field_name(headers)?; + let value = body.strip_suffix("\r\n").unwrap_or(body); + Some((name, value)) +} + +fn extract_multipart_field_name(headers: &str) -> Option<&str> { + let disposition = headers + .split("\r\n") + .find(|line| line.trim_start().to_ascii_lowercase().starts_with("content-disposition:"))?; + + let (_, attrs) = disposition.split_once(':')?; + for attr in attrs.split(';').skip(1) { + let (name, value) = attr.trim().split_once('=')?; + if !name.trim().eq_ignore_ascii_case("name") { + continue; + } + + let value = value.trim(); + let unquoted_double = value.strip_prefix('"').and_then(|inner| inner.strip_suffix('"')); + if let Some(result) = unquoted_double { + return Some(result); + } + + let unquoted_single = value.strip_prefix('\'').and_then(|inner| inner.strip_suffix('\'')); + if let Some(result) = unquoted_single { + return Some(result); + } + + return Some(value); + } + + None +} + +fn is_length_limit_error(err: &axum::Error) -> bool { + let mut current: Option<&(dyn StdError + 'static)> = Some(err); + while let Some(source) = current { + if source.is::() { + return true; + } + current = source.source(); + } + false +} + +fn is_multipart_content_type(content_type: &str) -> bool { + content_type + .split(';') + .next() + .is_some_and(|media_type| media_type.trim().eq_ignore_ascii_case("multipart/form-data")) +} + +fn extract_multipart_boundary(content_type: &str) -> Option<&str> { + for param in content_type.split(';').skip(1) { + let (name, value) = param.split_once('=')?; + if name.trim().eq_ignore_ascii_case("boundary") { + let trimmed = value.trim(); + let unquoted = trimmed + .strip_prefix('"') + .and_then(|inner| inner.strip_suffix('"')) + .unwrap_or(trimmed); + if !unquoted.is_empty() { + return Some(unquoted); + } + } + } + None +} + impl UserApiRequest { + fn set_field(&mut self, name: &str, value: &str) { + match name { + "username" => self.username = value.to_string(), + "password" => self.password = value.to_string(), + "token" => self.token = value.to_string(), + "action" => self.action = value.to_string(), + "series_id" => self.series_id = value.to_string(), + "vod_id" => self.vod_id = value.to_string(), + "stream_id" => self.stream_id = value.to_string(), + "category_id" => self.category_id = value.to_string(), + "limit" => self.limit = value.to_string(), + "start" => self.start = value.to_string(), + "end" => self.end = value.to_string(), + "stream" => self.stream = value.to_string(), + "duration" => self.duration = value.to_string(), + "type" | "content_type" => self.content_type = value.to_string(), + _ => {} + } + } + pub fn merge_prefer_primary(primary: &Self, fallback: &Self) -> Self { fn pick(primary: &str, fallback: &str) -> String { if primary.trim().is_empty() { @@ -97,7 +291,10 @@ impl UserApiRequest { #[cfg(test)] mod tests { - use super::UserApiRequest; + use super::{parse_body, UserApiRequest, MAX_BODY_SIZE_BYTES}; + use axum::body::Body; + use axum::extract::FromRequest; + use axum::http::{Request as HttpRequest, StatusCode}; #[test] fn merge_prefer_primary_uses_fallback_for_empty_fields() { @@ -194,4 +391,98 @@ mod tests { assert_eq!(merged.token, "query-token"); assert_eq!(merged.action, "query-action"); } + + #[tokio::test] + async fn parse_body_rejects_oversized_payloads() { + let oversized = "a".repeat(MAX_BODY_SIZE_BYTES + 1); + let err = parse_body(Body::from(oversized), "application/x-www-form-urlencoded") + .await + .expect_err("oversized bodies should be rejected"); + + match err { + super::ParseBodyError::PayloadTooLarge(msg) => { + assert!(msg.contains("body too large"), "unexpected error: {msg}"); + } + other => panic!("unexpected error: {other:?}"), + } + } + + #[tokio::test] + async fn extractor_surfaces_oversized_body_as_bad_request() { + let oversized = "a".repeat(MAX_BODY_SIZE_BYTES + 1); + let request = HttpRequest::builder() + .header("content-type", "application/x-www-form-urlencoded") + .uri("/player_api.php") + .body(Body::from(oversized)) + .expect("request should build"); + + let response = match super::UserApiRequestQueryOrBody::from_request(request, &()).await { + Ok(_) => panic!("oversized body should reject extractor"), + Err(response) => response, + }; + + assert_eq!(response.status(), StatusCode::PAYLOAD_TOO_LARGE); + } + + #[tokio::test] + async fn parse_body_detects_multipart_case_insensitively_and_preserves_field_whitespace() { + let body = concat!( + "--abc123\r\n", + "Content-Disposition: form-data; name=\"username\"\r\n\r\n", + " alice \r\n", + "--abc123\r\n", + "Content-Disposition: form-data; name=\"password\"\r\n\r\n", + "\t secret \t\r\n", + "--abc123--\r\n", + ); + + let parsed = parse_body( + Body::from(body), + "Multipart/Form-Data; charset=utf-8; boundary=\"abc123\"", + ) + .await + .expect("multipart body should parse"); + + assert_eq!(parsed.username, " alice "); + assert_eq!(parsed.password, "\t secret \t"); + } + + #[tokio::test] + async fn parse_body_accepts_content_type_field_name_in_multipart() { + let body = concat!( + "--abc123\r\n", + "Content-Disposition: form-data; name=\"content_type\"\r\n\r\n", + "m3u_plus\r\n", + "--abc123--\r\n", + ); + + let parsed = parse_body( + Body::from(body), + "multipart/form-data; boundary=abc123", + ) + .await + .expect("multipart body should parse"); + + assert_eq!(parsed.content_type, "m3u_plus"); + } + + #[tokio::test] + async fn parse_body_uses_content_disposition_name_not_filename() { + let body = concat!( + "--abc123\r\n", + "Content-Disposition: form-data; filename=\"name=\\\"wrong\\\".txt\"; name=\"username\"\r\n\r\n", + "alice\r\n", + "--abc123--\r\n", + ); + + let parsed = parse_body( + Body::from(body), + "multipart/form-data; boundary=abc123", + ) + .await + .expect("multipart body should parse"); + + assert_eq!(parsed.username, "alice"); + } + } diff --git a/backend/src/api/model/streams/active_client_stream.rs b/backend/src/api/model/streams/active_client_stream.rs index 0f7ec48bc..0b80f4944 100644 --- a/backend/src/api/model/streams/active_client_stream.rs +++ b/backend/src/api/model/streams/active_client_stream.rs @@ -1044,7 +1044,7 @@ mod tests { StreamError, UpdateGuard, }, auth::Fingerprint, - model::{AppConfig, Config, ConfigInput, GracePeriodOptions, ProcessTargets, ProxyUserCredentials, SourcesConfig}, + model::{AppConfig, Config, ConfigInput, GracePeriodOptions, MediaToolCapabilities, ProcessTargets, ProxyUserCredentials, SourcesConfig}, utils::{FileLockManager, GeoIp}, }; use arc_swap::{ArcSwap, ArcSwapOption}; @@ -1106,7 +1106,7 @@ mod tests { custom_stream_response: Arc::new(ArcSwapOption::default()), access_token_secret: [0; 32], encrypt_secret: [0; 16], - ffprobe_available: Arc::default(), + media_tools: Arc::new(MediaToolCapabilities::new()), } } diff --git a/backend/src/api/model/streams/provider_stream.rs b/backend/src/api/model/streams/provider_stream.rs index 1965d0040..f5fe921a7 100644 --- a/backend/src/api/model/streams/provider_stream.rs +++ b/backend/src/api/model/streams/provider_stream.rs @@ -287,7 +287,7 @@ mod tests { use super::{create_channel_unavailable_stream, CustomVideoStreamType}; use crate::{ api::model::TransportStreamBuffer, - model::{AppConfig, Config, ConfigInput, CustomStreamResponse, SourcesConfig}, + model::{AppConfig, Config, ConfigInput, CustomStreamResponse, MediaToolCapabilities, SourcesConfig}, utils::FileLockManager, }; use arc_swap::{ArcSwap, ArcSwapOption}; @@ -338,7 +338,7 @@ mod tests { custom_stream_response: Arc::new(ArcSwapOption::default()), access_token_secret: [0; 32], encrypt_secret: [0; 16], - ffprobe_available: Arc::default(), + media_tools: Arc::new(MediaToolCapabilities::new()), }; let mut ts_packet = vec![0_u8; 188]; diff --git a/backend/src/api/model/streams/shared_stream_manager.rs b/backend/src/api/model/streams/shared_stream_manager.rs index 47a4933ae..3399f682d 100644 --- a/backend/src/api/model/streams/shared_stream_manager.rs +++ b/backend/src/api/model/streams/shared_stream_manager.rs @@ -749,7 +749,7 @@ mod tests { use super::{SharedStreamManager, SharedStreamState, CHANNEL_SIZE}; use crate::{ api::model::{ActiveProviderManager, EventManager}, - model::{AppConfig, Config, ConfigInput, SourcesConfig}, + model::{AppConfig, Config, ConfigInput, MediaToolCapabilities, SourcesConfig}, utils::FileLockManager, }; use arc_swap::{ArcSwap, ArcSwapOption}; @@ -803,7 +803,7 @@ mod tests { custom_stream_response: Arc::new(ArcSwapOption::default()), access_token_secret: [0; 32], encrypt_secret: [0; 16], - ffprobe_available: Arc::default(), + media_tools: Arc::new(MediaToolCapabilities::new()), } } diff --git a/backend/src/library/processor.rs b/backend/src/library/processor.rs index 436d812b5..5da5ecf4a 100644 --- a/backend/src/library/processor.rs +++ b/backend/src/library/processor.rs @@ -5,6 +5,7 @@ use crate::library::metadata_storage::MetadataStorage; use crate::library::scanner::LibraryScanner; use crate::library::{MediaGroup, MediaGrouper, thumbnail::{self, ThumbnailExtractor}}; use crate::model::{AppConfig, LibraryConfig, MetadataUpdateConfig}; +use crate::utils::ffmpeg::FfmpegExecutor; use log::{debug, error, info, warn}; use path_clean::PathClean; use shared::model::{LibraryMetadataFormat, LibraryScanResult}; @@ -135,9 +136,23 @@ impl LibraryProcessor { false }; + let ffmpeg_available = if self.thumbnail_extractor.is_some() { + if let Some(app_cfg) = &self.app_config { + app_cfg.is_ffmpeg_available().await + } else { + FfmpegExecutor::new().check_ffmpeg_availability().await + } + } else { + false + }; + + if self.thumbnail_extractor.is_some() && !ffmpeg_available { + warn!("Thumbnail extraction disabled because ffmpeg is unavailable"); + } + // Process each scanned file for group in &media_groups { - match self.process_group(group, &existing_map, force_rescan, ffprobe_enabled).await { + match self.process_group(group, &existing_map, force_rescan, ffprobe_enabled, ffmpeg_available).await { Ok(action) => match action { ProcessAction::Added => result.files_added += 1, ProcessAction::Updated => result.files_updated += 1, @@ -170,7 +185,7 @@ impl LibraryProcessor { } } - if self.thumbnail_extractor.is_some() { + if ffmpeg_available { self.storage.cleanup_orphaned_thumbnails().await; } @@ -178,13 +193,20 @@ impl LibraryProcessor { Ok(result) } - async fn process_group(&self, group: &MediaGroup, existing_map: &HashMap, force_rescan: bool, can_probe: bool) -> Result { + async fn process_group( + &self, + group: &MediaGroup, + existing_map: &HashMap, + force_rescan: bool, + can_probe: bool, + can_extract_thumbnails: bool, + ) -> Result { match group { MediaGroup::Movie { file: _, .. } => { - self.process_movie(group, existing_map, force_rescan, can_probe).await + self.process_movie(group, existing_map, force_rescan, can_probe, can_extract_thumbnails).await } MediaGroup::Series { show_key: _, episodes: _ } => { - self.process_series_group(group, existing_map, force_rescan, can_probe).await + self.process_series_group(group, existing_map, force_rescan, can_probe, can_extract_thumbnails).await } } } @@ -210,7 +232,14 @@ impl LibraryProcessor { //} // Processes a single video file - async fn process_movie(&self, group: &MediaGroup, existing_map: &HashMap, force_rescan: bool, _can_probe: bool) -> Result { + async fn process_movie( + &self, + group: &MediaGroup, + existing_map: &HashMap, + force_rescan: bool, + _can_probe: bool, + can_extract_thumbnails: bool, + ) -> Result { let MediaGroup::Movie { file, .. } = group else { return Err(format!("Expected movie to resolve but got {group}")) }; // Check if file already exists in cache let (mut cache_entry, status) = if let Some(existing_entry) = existing_map.get(&file.file_path) { @@ -251,7 +280,12 @@ impl LibraryProcessor { (entry, ProcessAction::Added) }; - self.extract_thumbnail_if_needed(&mut cache_entry, &file.file_path, file.modified_timestamp).await; + self.extract_thumbnail_if_needed( + &mut cache_entry, + &file.file_path, + file.modified_timestamp, + can_extract_thumbnails, + ).await; self.storage.store(&cache_entry).await.map_err(|e| e.to_string())?; self.write_metadata_files(&cache_entry).await.map_err(|e| e.to_string())?; Ok(status) @@ -263,7 +297,9 @@ impl LibraryProcessor { &self, group: &MediaGroup, existing_map: &HashMap, - force_rescan: bool, _can_probe: bool + force_rescan: bool, + _can_probe: bool, + can_extract_thumbnails: bool, ) -> Result { let MediaGroup::Series { show_key, episodes } = group else { return Err(format!("Expected series to resolve but got {group}")) }; let series_file_path = episodes @@ -370,6 +406,7 @@ impl LibraryProcessor { &episode.file.file_path, episode.file.modified_timestamp, prev_mtime, + can_extract_thumbnails, ).await; } else { let mut new_episode = series_episode.clone(); @@ -381,6 +418,7 @@ impl LibraryProcessor { &episode.file.file_path, episode.file.modified_timestamp, Some(previous_file_modified), + can_extract_thumbnails, ).await; double_episodes.push(new_episode); } @@ -404,6 +442,7 @@ impl LibraryProcessor { &mut chache_entry, &first_ep.file.file_path, first_ep.file.modified_timestamp, + can_extract_thumbnails, ).await; } } @@ -426,7 +465,12 @@ impl LibraryProcessor { cache_entry: &mut MetadataCacheEntry, file_path: &str, file_mtime: i64, + can_extract_thumbnails: bool, ) { + if !can_extract_thumbnails { + return; + } + // Skip if already has a poster from TMDB/NFO if cache_entry.metadata.poster().is_some() { // Clear stale generated-thumbnail references so they can be reclaimed @@ -471,6 +515,7 @@ impl LibraryProcessor { file_path: &str, file_mtime: i64, previous_file_mtime: Option, + can_extract_thumbnails: bool, ) { if !episode.thumb.as_deref().unwrap_or_default().is_empty() && episode.thumbnail_id.as_deref().unwrap_or_default().is_empty() @@ -478,7 +523,12 @@ impl LibraryProcessor { return; } - if let Some(thumbnail_id) = self.extract_thumbnail_id_for_file(file_path, file_mtime, previous_file_mtime).await { + if let Some(thumbnail_id) = self.extract_thumbnail_id_for_file( + file_path, + file_mtime, + previous_file_mtime, + can_extract_thumbnails, + ).await { episode.thumbnail_id = Some(thumbnail_id); } } @@ -488,7 +538,12 @@ impl LibraryProcessor { file_path: &str, file_mtime: i64, previous_file_mtime: Option, + can_extract_thumbnails: bool, ) -> Option { + if !can_extract_thumbnails { + return None; + } + let extractor = self.thumbnail_extractor.as_ref()?; let hash = thumbnail::file_hash(file_path); diff --git a/backend/src/library/thumbnail.rs b/backend/src/library/thumbnail.rs index 147657e56..c0442d4d2 100644 --- a/backend/src/library/thumbnail.rs +++ b/backend/src/library/thumbnail.rs @@ -1,10 +1,6 @@ use crate::model::ThumbnailConfig; -use log::debug; +use crate::utils::ffmpeg::FfmpegExecutor; use std::path::Path; -use std::time::Duration; -use tokio::process::Command; - -const FFMPEG_TIMEOUT: Duration = Duration::from_secs(60); /// Computes a stable BLAKE3 hash for a file path (or URL). pub fn file_hash(path: &str) -> String { @@ -13,16 +9,14 @@ pub fn file_hash(path: &str) -> String { pub struct ThumbnailExtractor { config: ThumbnailConfig, + ffmpeg: FfmpegExecutor, } impl ThumbnailExtractor { - pub fn new(config: ThumbnailConfig) -> Self { Self { config } } - - /// Checks if the system `ffmpeg` binary is available. - pub async fn check_ffmpeg_availability() -> bool { - match Command::new("ffmpeg").arg("-version").output().await { - Ok(output) => output.status.success(), - Err(_) => false, + pub fn new(config: ThumbnailConfig) -> Self { + Self { + config, + ffmpeg: FfmpegExecutor::new(), } } @@ -37,85 +31,14 @@ impl ThumbnailExtractor { return Err(format!("File not found: {file_path}")); } - self.run_ffmpeg(file_path).await + self.ffmpeg + .create_thumbnail(file_path, self.config.width, self.config.height) + .await } // TODO: Future feature: support thumbnail extraction from remote HTTP(S) // inputs via ranged reads without tying this logic into the local library // scan path yet. - - fn build_scale_filter(&self) -> String { - format!( - "scale={}:{}:force_original_aspect_ratio=increase,crop={}:{}", - self.config.width, self.config.height, self.config.width, self.config.height - ) - } - - /// Runs ffmpeg to extract a single frame at ~180s into the video. - /// Falls back to position 0 if the video is shorter than 180s. - async fn run_ffmpeg(&self, input_path: &str) -> Result, String> { - let temp_dir = tempfile::tempdir() - .map_err(|e| format!("Failed to create temp dir: {e}"))?; - let output_path = temp_dir.path().join("thumb.jpg"); - let scale_filter = self.build_scale_filter(); - - let output = self.run_ffmpeg_with_timeout(&[ - "-ss", "180", - "-i", input_path, - "-frames:v", "1", - "-vf", &scale_filter, - "-q:v", "1", - "-y", - &output_path.to_string_lossy(), - ]).await - .map_err(|e| format!("Failed to run ffmpeg: {e}"))?; - - if !output.status.success() { - let stderr = String::from_utf8_lossy(&output.stderr); - // Retry at position 0 if seeking past end of short video - if stderr.contains("Output file is empty") || stderr.contains("nothing was encoded") { - debug!("Video shorter than 180s, retrying at position 0: {input_path}"); - let output = self.run_ffmpeg_with_timeout(&[ - "-ss", "0", - "-i", input_path, - "-frames:v", "1", - "-vf", &scale_filter, - "-q:v", "1", - "-y", - &output_path.to_string_lossy(), - ]).await - .map_err(|e| format!("Failed to run ffmpeg retry: {e}"))?; - - if !output.status.success() { - let stderr = String::from_utf8_lossy(&output.stderr); - return Err(format!("ffmpeg failed at position 0: {stderr}")); - } - } else { - return Err(format!("ffmpeg failed: {stderr}")); - } - } - - tokio::fs::read(&output_path) - .await - .map_err(|e| format!("Failed to read thumbnail: {e}")) - } - - /// Spawns ffmpeg with the given args and enforces a bounded timeout. - /// On timeout the child process is killed via `kill_on_drop`. - async fn run_ffmpeg_with_timeout(&self, args: &[&str]) -> Result { - let child = Command::new("ffmpeg") - .args(args) - .stdout(std::process::Stdio::piped()) - .stderr(std::process::Stdio::piped()) - .kill_on_drop(true) - .spawn() - .map_err(|e| format!("Failed to spawn ffmpeg: {e}"))?; - - tokio::time::timeout(FFMPEG_TIMEOUT, child.wait_with_output()) - .await - .map_err(|_| "Timed out running ffmpeg".to_string())? - .map_err(|e| e.to_string()) - } } #[cfg(test)] @@ -141,4 +64,4 @@ mod tests { let url = format!("/api/v1/library/thumbnail/{}", "test-uuid-123"); assert_eq!(url, "/api/v1/library/thumbnail/test-uuid-123"); } -} \ No newline at end of file +} diff --git a/backend/src/messaging.rs b/backend/src/messaging.rs index 81fb876a2..414918fde 100644 --- a/backend/src/messaging.rs +++ b/backend/src/messaging.rs @@ -296,7 +296,7 @@ async fn resolve_template<'a>(app_config: &'a Arc, http_client: &'a r #[cfg(test)] mod tests { use arc_swap::{ArcSwap, ArcSwapOption}; - use crate::model::ProcessingStats; + use crate::model::{MediaToolCapabilities, ProcessingStats}; use super::*; use shared::model::{ConfigPaths}; use crate::utils::FileLockManager; @@ -324,7 +324,7 @@ mod tests { custom_stream_response: Arc::new(ArcSwapOption::default()), access_token_secret: [1,2,3,4,5,6,7,8,9,10,11,12,13,14,15,16,17,18,19,20,21,22,23,24,25,26,27,28,29,30,31,32], encrypt_secret: [1,2,3,4,5,6,7,8,9,10,11,12,13,14,15,16], - ffprobe_available: Arc::default(), + media_tools: Arc::new(MediaToolCapabilities::new()), }) } diff --git a/backend/src/model/config/app.rs b/backend/src/model/config/app.rs index 3445a8914..f4b2e9498 100644 --- a/backend/src/model/config/app.rs +++ b/backend/src/model/config/app.rs @@ -1,11 +1,10 @@ use crate::api::model::TransportStreamBuffer; use crate::model::{ ApiProxyConfig, ApiProxyServerInfo, Config, ConfigInput, ConfigInputOptions, ConfigTarget, CustomStreamResponse, - GracePeriodOptions, HdHomeRunConfig, HdHomeRunFlags, Mappings, ProxyUserCredentials, + GracePeriodOptions, HdHomeRunConfig, HdHomeRunFlags, Mappings, MediaToolCapabilities, ProxyUserCredentials, ReverseProxyDisabledHeaderConfig, SourcesConfig, TargetOutput, }; use crate::utils; -use crate::utils::ffmpeg::check_ffprobe_availability; use arc_swap::{ArcSwap, ArcSwapOption}; use log::{error, warn}; use rand::Rng; @@ -22,7 +21,6 @@ use std::fs::File; use std::io::Read; use std::path::{Path, PathBuf}; use std::sync::Arc; -use tokio::sync::OnceCell; fn generate_secret() -> [u8; 32] { let mut rng = rand::rng(); @@ -42,7 +40,7 @@ pub struct AppConfig { pub custom_stream_response: Arc>, pub access_token_secret: [u8; 32], pub encrypt_secret: [u8; 16], - pub(crate) ffprobe_available: Arc>, + pub(crate) media_tools: Arc, } impl AppConfig { @@ -481,8 +479,10 @@ impl AppConfig { return false; } - *self.ffprobe_available.get_or_init(|| async { - check_ffprobe_availability().await - }).await + self.media_tools.is_ffprobe_available().await + } + + pub async fn is_ffmpeg_available(&self) -> bool { + self.media_tools.is_ffmpeg_available().await } } diff --git a/backend/src/model/config/media_tools.rs b/backend/src/model/config/media_tools.rs new file mode 100644 index 000000000..486f579d5 --- /dev/null +++ b/backend/src/model/config/media_tools.rs @@ -0,0 +1,48 @@ +use crate::utils::ffmpeg::FfmpegExecutor; +use shared::create_bitset; +use tokio::sync::OnceCell; + +create_bitset!(u8, MediaToolCapability, Ffmpeg, Ffprobe); + +#[derive(Debug, Default)] +pub struct MediaToolCapabilities { + available: OnceCell, +} + +impl MediaToolCapabilities { + #[must_use] + pub fn new() -> Self { + Self { + available: OnceCell::new(), + } + } + + pub async fn is_ffmpeg_available(&self) -> bool { + self.available().await.contains(MediaToolCapability::Ffmpeg) + } + + pub async fn is_ffprobe_available(&self) -> bool { + self.available().await.contains(MediaToolCapability::Ffprobe) + } + + async fn available(&self) -> MediaToolCapabilitySet { + *self.available.get_or_init(Self::detect_available_tools).await + } + + async fn detect_available_tools() -> MediaToolCapabilitySet { + let executor = FfmpegExecutor::new(); + let ffmpeg = executor.check_ffmpeg_availability(); + let ffprobe = executor.check_ffprobe_availability(); + let (ffmpeg_available, ffprobe_available) = tokio::join!(ffmpeg, ffprobe); + + let mut available = MediaToolCapabilitySet::new(); + if ffmpeg_available { + available.set(MediaToolCapability::Ffmpeg); + } + if ffprobe_available { + available.set(MediaToolCapability::Ffprobe); + } + available + } +} + diff --git a/backend/src/model/config/mod.rs b/backend/src/model/config/mod.rs index 917cd4b96..9b967f2f2 100644 --- a/backend/src/model/config/mod.rs +++ b/backend/src/model/config/mod.rs @@ -4,6 +4,7 @@ mod web_ui; mod web_auth; mod messaging; mod metadata_update; +mod media_tools; mod hdhomerun; mod ip_check; mod source; @@ -46,6 +47,7 @@ pub use input::*; pub use ip_check::*; pub use log::*; pub use messaging::*; +pub use media_tools::*; pub use metadata_update::*; pub use proxy::*; pub use rate_limit::*; diff --git a/backend/src/modules.rs b/backend/src/modules.rs index c19479738..fd45858cb 100644 --- a/backend/src/modules.rs +++ b/backend/src/modules.rs @@ -12,7 +12,6 @@ macro_rules! include_modules { pub mod processing; pub mod ptt; pub mod repository; - pub mod tools; pub mod utils; }; } diff --git a/backend/src/processing/parser/xtream.rs b/backend/src/processing/parser/xtream.rs index d33f02f99..38d5ddc6a 100644 --- a/backend/src/processing/parser/xtream.rs +++ b/backend/src/processing/parser/xtream.rs @@ -11,9 +11,29 @@ use shared::model::{EpisodeStreamProperties, LiveStreamProperties, PlaylistGroup SeriesStreamProperties, StreamProperties, VideoStreamProperties, XtreamCluster, XtreamPlaylistItem}; use shared::utils::{generate_playlist_uuid, trim_last_slash, Internable}; +use std::collections::HashMap; use std::sync::Arc; use tokio::task::spawn_blocking; +/// Bucket size for composite ordinal encoding in the streaming parser. +/// Layout: `cat_position * CAT_BUCKET + within_cat_counter`. +/// Supports up to ~42 900 categories with up to 100 000 streams each. +const CAT_BUCKET: u32 = 100_000; + +fn next_group_map_key(group_map: &IndexMap) -> u32 { + if let Some(max_key) = group_map.keys().copied().max() { + if let Some(next_key) = max_key.checked_add(1) { + return next_key; + } + } + + let mut candidate = 0u32; + while group_map.contains_key(&candidate) { + candidate = candidate.saturating_add(1); + } + candidate +} + async fn map_to_xtream_category(categories: DynReader, input_name: &Arc) -> Result, TuliproxError> { let input_name_clone = Arc::clone(input_name); spawn_blocking(move || { @@ -174,12 +194,12 @@ pub async fn parse_xtream(input: &ConfigInput, input.has_flag(ConfigInputFlags::XtreamLiveStreamWithoutExtension), ); - for (ord_counter, stream) in (1_u32..).zip(xtream_streams) { + for stream in xtream_streams { let group = group_map.get_mut(&stream.get_category_id()).unwrap_or(&mut unknown_grp); let category_name = &group.category_name; let stream_url = create_xtream_url(xtream_cluster, url, username, password, &stream, live_stream_use_prefix, live_stream_without_extension); let item_type = PlaylistItemType::from(xtream_cluster); - let mut item = PlaylistItem { + let item = PlaylistItem { header: PlaylistItemHeader { id: stream.get_stream_id().intern(), uuid: generate_playlist_uuid(&input_name, &stream.get_stream_id().to_string(), item_type, &stream_url), @@ -197,14 +217,25 @@ pub async fn parse_xtream(input: &ConfigInput, ..Default::default() }, }; - item.header.source_ordinal = ord_counter; group.add(item); } + if !unknown_grp.channels.is_empty() { + let unknown_key = next_group_map_key(&group_map); + unknown_grp.category_id = unknown_key; + group_map.insert(unknown_key, unknown_grp); + } - let has_channels = !unknown_grp.channels.is_empty(); - if has_channels { - group_map.insert(0, unknown_grp); + // Assign source_ordinal in category-list order (primary) + // with stream-list position within each category (secondary). + // The IndexMap preserves the provider's category ordering, + // so a single sequential pass produces the correct ordinals. + let mut ordinal: u32 = 0; + for category in group_map.values_mut() { + for channel in &mut category.channels { + ordinal += 1; + channel.header.source_ordinal = ordinal; + } } Ok(Some(group_map.values().filter(|category| !category.channels.is_empty()) @@ -251,42 +282,62 @@ where let group_map: IndexMap> = xtream_categories.iter().map(|c| (c.category_id, c.category_name.clone())).collect(); let unknown_group_name = "Unknown".intern(); + // Category position lookup for source_ordinal: streams are ordered by + // category-list position (primary) then arrival order within that + // category (secondary). We encode both into a single u32 so the + // downstream sort-by-source_ordinal reproduces the provider's category + // ordering without any extra fields or a second sort pass. + // + // Layout: cat_position * CAT_BUCKET + within_cat_counter + // With CAT_BUCKET = 100_000 this supports up to ~42_900 categories + // with up to 100 000 streams each — well beyond real-world sizes. + let cat_order: HashMap = group_map.keys().enumerate() + .map(|(idx, &cat_id)| (cat_id, u32::try_from(idx).unwrap_or(u32::MAX))) + .collect(); + let unknown_cat_pos = u32::try_from(cat_order.len()).unwrap_or(u32::MAX); + spawn_blocking(move || { let reader = tokio_util::io::SyncIoBridge::new(streams); let mut deserializer = serde_json::Deserializer::from_reader(reader); - let mut source_ordinal = 0u32; + let mut cat_counters: HashMap = HashMap::new(); + let mut cat_source_ordinal = |cat_id: u32| -> u32 { + let cat_pos = cat_order.get(&cat_id).copied().unwrap_or(unknown_cat_pos); + let counter = cat_counters.entry(cat_pos).or_insert(0); + *counter += 1; + cat_pos.saturating_mul(CAT_BUCKET).saturating_add(*counter) + }; match xtream_cluster { XtreamCluster::Live => { let mut on_stream = |stream: LiveStreamProperties| { - source_ordinal += 1; + let ordinal = cat_source_ordinal(stream.category_id); let stream_prop = StreamProperties::Live(Box::new(stream)); process_stream_item(&input_name, &url, &username, &password, xtream_cluster, &group_map, &unknown_group_name, - stream_prop, &mut on_item, live_stream_use_prefix, live_stream_without_extension, source_ordinal) + stream_prop, &mut on_item, live_stream_use_prefix, live_stream_without_extension, ordinal) }; let visitor = XtreamItemVisitor { on_item: &mut on_stream, _marker: std::marker::PhantomData }; deserializer.deserialize_any(visitor).map_err(|e| notify_err!("JSON parse error: {e}"))?; } XtreamCluster::Video => { let mut on_stream = |stream: VideoStreamProperties| { - source_ordinal += 1; + let ordinal = cat_source_ordinal(stream.category_id); let stream_prop = StreamProperties::Video(Box::new(stream)); process_stream_item(&input_name, &url, &username, &password, xtream_cluster, &group_map, &unknown_group_name, - stream_prop, &mut on_item, live_stream_use_prefix, live_stream_without_extension, source_ordinal) + stream_prop, &mut on_item, live_stream_use_prefix, live_stream_without_extension, ordinal) }; let visitor = XtreamItemVisitor { on_item: &mut on_stream, _marker: std::marker::PhantomData }; deserializer.deserialize_any(visitor).map_err(|e| notify_err!("JSON parse error: {e}"))?; } XtreamCluster::Series => { let mut on_stream = |stream: SeriesStreamProperties| { - source_ordinal += 1; + let ordinal = cat_source_ordinal(stream.category_id); let stream_prop = StreamProperties::Series(Box::new(stream)); process_stream_item(&input_name, &url, &username, &password, xtream_cluster, &group_map, &unknown_group_name, - stream_prop, &mut on_item, live_stream_use_prefix, live_stream_without_extension, source_ordinal) + stream_prop, &mut on_item, live_stream_use_prefix, live_stream_without_extension, ordinal) }; let visitor = XtreamItemVisitor { on_item: &mut on_stream, _marker: std::marker::PhantomData }; deserializer.deserialize_any(visitor).map_err(|e| notify_err!("JSON parse error: {e}"))?; @@ -390,17 +441,40 @@ where #[cfg(test)] mod tests { - use super::parse_xtream_series_info; + use super::CAT_BUCKET; + use super::{parse_xtream, parse_xtream_series_info, parse_xtream_streaming}; use crate::processing::parser::xtream::map_to_xtream_streams; use crate::model::ConfigInput; - use crate::utils::async_file_reader; + use crate::utils::{async_file_reader, request::DynReader}; use shared::model::{ UUIDType, SeriesStreamDetailEpisodeProperties, SeriesStreamDetailProperties, SeriesStreamProperties, - XtreamCluster, XtreamSeriesInfo, + XtreamCluster, XtreamPlaylistItem, XtreamSeriesInfo, }; use shared::utils::Internable; use std::fs; + use std::sync::{Arc, Mutex}; + use tokio::io::AsyncWriteExt; + + fn make_reader(content: &str) -> DynReader { + let (mut writer, reader) = tokio::io::duplex(4096); + let bytes = content.as_bytes().to_vec(); + tokio::spawn(async move { + writer.write_all(&bytes).await.unwrap(); + writer.shutdown().await.unwrap(); + }); + Box::pin(reader) + } + + fn test_input() -> ConfigInput { + ConfigInput { + name: "input".intern(), + url: "http://provider.example".to_string(), + username: Some("user".to_string()), + password: Some("pass".to_string()), + ..ConfigInput::default() + } + } #[test] fn test_read_json_file_into_struct() { @@ -497,13 +571,7 @@ mod tests { #[test] fn test_parse_xtream_series_info_keeps_zero_parent_source_ordinal() { - let input = ConfigInput { - name: "input".intern(), - url: "http://provider.example".to_string(), - username: Some("user".to_string()), - password: Some("pass".to_string()), - ..ConfigInput::default() - }; + let input = test_input(); let episode: SeriesStreamDetailEpisodeProperties = serde_json::from_str( r#"{"id":101,"episode_num":1,"season":1,"title":"S01E01","container_extension":"mp4"}"#, @@ -533,4 +601,126 @@ mod tests { assert_eq!(episodes[0].header.source_ordinal, 0); } + + #[tokio::test] + async fn test_parse_xtream_streaming_source_ordinal_uses_category_then_channel_position() { + let categories = r#" + [ + {"category_id":"20","category_name":"News"}, + {"category_id":"10","category_name":"Sports"} + ] + "#; + let streams = r#" + [ + {"name":"sports-1","stream_id":101,"category_id":"10","added":"0"}, + {"name":"news-1","stream_id":201,"category_id":"20","added":"0"}, + {"name":"sports-2","stream_id":102,"category_id":"10","added":"0"}, + {"name":"news-2","stream_id":202,"category_id":"20","added":"0"} + ] + "#; + + let items: Arc>> = Arc::new(Mutex::new(Vec::new())); + let sink = Arc::clone(&items); + + parse_xtream_streaming( + &test_input(), + XtreamCluster::Live, + make_reader(categories), + make_reader(streams), + move |item| { + sink.lock().unwrap().push(item); + Ok(()) + }, + ) + .await + .unwrap(); + + let mut items = items.lock().unwrap().clone(); + items.sort_by_key(|item| item.source_ordinal); + + let names: Vec<&str> = items.iter().map(|item| item.name.as_ref()).collect(); + assert_eq!(names, vec!["news-1", "news-2", "sports-1", "sports-2"]); + assert_eq!(items[0].source_ordinal, 1); + assert_eq!(items[1].source_ordinal, 2); + assert_eq!(items[2].source_ordinal, CAT_BUCKET + 1); + assert_eq!(items[3].source_ordinal, CAT_BUCKET + 2); + } + + #[tokio::test] + async fn test_parse_xtream_groups_multiple_unknown_categories_into_unknown_group() { + let categories = r#" + [ + {"category_id":"20","category_name":"Known"} + ] + "#; + let streams = r#" + [ + {"name":"unknown-999","stream_id":301,"category_id":"999","added":"0"}, + {"name":"known-1","stream_id":201,"category_id":"20","added":"0"}, + {"name":"unknown-888","stream_id":302,"category_id":"888","added":"0"} + ] + "#; + + let groups = parse_xtream( + &test_input(), + XtreamCluster::Live, + make_reader(categories), + make_reader(streams), + ) + .await + .unwrap() + .unwrap(); + + assert_eq!(groups.len(), 2); + assert_eq!(groups[0].title.as_ref(), "Known"); + assert_eq!(groups[0].channels.len(), 1); + assert_eq!(groups[0].channels[0].header.name.as_ref(), "known-1"); + assert_eq!(groups[0].channels[0].header.source_ordinal, 1); + + assert_eq!(groups[1].title.as_ref(), "Unknown"); + let unknown_names: Vec<&str> = groups[1] + .channels + .iter() + .map(|item| item.header.name.as_ref()) + .collect(); + assert_eq!(unknown_names, vec!["unknown-999", "unknown-888"]); + assert_eq!(groups[1].channels[0].header.source_ordinal, 2); + assert_eq!(groups[1].channels[1].header.source_ordinal, 3); + } + + #[tokio::test] + async fn test_parse_xtream_keeps_real_category_zero_and_appends_unknown_group() { + let categories = r#" + [ + {"category_id":"0","category_name":"Provider Zero"} + ] + "#; + let streams = r#" + [ + {"name":"zero-1","stream_id":201,"category_id":"0","added":"0"}, + {"name":"unknown-1","stream_id":301,"category_id":"999","added":"0"} + ] + "#; + + let groups = parse_xtream( + &test_input(), + XtreamCluster::Live, + make_reader(categories), + make_reader(streams), + ) + .await + .unwrap() + .unwrap(); + + assert_eq!(groups.len(), 2); + assert_eq!(groups[0].id, 0); + assert_eq!(groups[0].title.as_ref(), "Provider Zero"); + assert_eq!(groups[0].channels[0].header.name.as_ref(), "zero-1"); + assert_eq!(groups[0].channels[0].header.source_ordinal, 1); + + assert_eq!(groups[1].title.as_ref(), "Unknown"); + assert_ne!(groups[1].id, 0); + assert_eq!(groups[1].channels[0].header.name.as_ref(), "unknown-1"); + assert_eq!(groups[1].channels[0].header.source_ordinal, 2); + } } diff --git a/backend/src/processing/processor/probe_handle_guard.rs b/backend/src/processing/processor/probe_handle_guard.rs index 712675f27..a0b3b7cca 100644 --- a/backend/src/processing/processor/probe_handle_guard.rs +++ b/backend/src/processing/processor/probe_handle_guard.rs @@ -43,7 +43,7 @@ mod tests { use super::ProbeHandleGuard; use crate::{ api::model::{ActiveProviderManager, EventManager}, - model::{AppConfig, Config, ConfigInput, SourcesConfig}, + model::{AppConfig, Config, ConfigInput, MediaToolCapabilities, SourcesConfig}, utils::FileLockManager, }; use arc_swap::{ArcSwap, ArcSwapOption}; @@ -93,7 +93,7 @@ mod tests { custom_stream_response: Arc::new(ArcSwapOption::default()), access_token_secret: [0; 32], encrypt_secret: [0; 16], - ffprobe_available: Arc::default(), + media_tools: Arc::new(MediaToolCapabilities::new()), } } diff --git a/backend/src/processing/processor/stream_probe.rs b/backend/src/processing/processor/stream_probe.rs index fe2fe3b15..5df357968 100644 --- a/backend/src/processing/processor/stream_probe.rs +++ b/backend/src/processing/processor/stream_probe.rs @@ -3,8 +3,8 @@ use crate::model::ConfigInput; use crate::model::{AppConfig}; use crate::processing::processor::{select_cancel_token, ProbeHandleGuard}; use crate::repository::{get_input_m3u_playlist_file_path, get_input_storage_path, get_input_local_library_playlist_file_path, xtream_get_file_path, BPlusTreeUpdate}; -use crate::utils::{debug_if_enabled, ffmpeg}; -use crate::utils::ffmpeg::{ProbeFailureKind, ProbeUrlOutcome}; +use crate::utils::debug_if_enabled; +use crate::utils::ffmpeg::{FfmpegExecutor, ProbeFailureKind, ProbeUrlOutcome}; use log::{info, warn}; use shared::error::TuliproxError; use shared::model::{EpisodeStreamProperties, InputType, PlaylistItemType, StreamProperties, VideoStreamDetailProperties, VideoStreamProperties, LiveStreamProperties, M3uPlaylistItem, XtreamCluster, XtreamPlaylistItem}; @@ -124,7 +124,7 @@ pub async fn update_generic_stream_metadata( acquired_handle.as_ref().and_then(ProbeHandleGuard::handle), active_handle, ); - let probe_data = ffmpeg::probe_url_with_cancel( + let probe_data = FfmpegExecutor::new().probe_url_with_cancel( &probe_url, user_agent.as_deref(), analyze_duration, diff --git a/backend/src/processing/processor/xtream.rs b/backend/src/processing/processor/xtream.rs index 5779d8949..08f7e409f 100644 --- a/backend/src/processing/processor/xtream.rs +++ b/backend/src/processing/processor/xtream.rs @@ -5,7 +5,7 @@ use crate::model::{AppConfig, ConfigInput, ConfigInputFlags}; use shared::model::{LiveStreamProperties, StreamProperties, XtreamCluster, XtreamPlaylistItem}; use crate::repository::{get_input_storage_path, persist_input_live_info, BPlusTreeQuery, xtream_get_file_path}; use crate::utils::{debug_if_enabled}; -use crate::utils::ffmpeg::{ProbeFailureKind, ProbeUrlOutcome}; +use crate::utils::ffmpeg::{FfmpegExecutor, ProbeFailureKind, ProbeUrlOutcome}; use log::{debug, warn}; use crate::processing::parser::xtream::create_xtream_url; use crate::api::model::{ActiveProviderManager, ProviderHandle, ProviderIdType}; @@ -133,7 +133,7 @@ pub async fn update_live_stream_metadata( let mut success = false; let mut not_found = false; - match crate::utils::ffmpeg::probe_url( + match FfmpegExecutor::new().probe_url( &stream_url, user_agent.as_deref(), analyze_duration, diff --git a/backend/src/processing/processor/xtream_series.rs b/backend/src/processing/processor/xtream_series.rs index 5aa7da0bf..abb3aac6c 100644 --- a/backend/src/processing/processor/xtream_series.rs +++ b/backend/src/processing/processor/xtream_series.rs @@ -19,7 +19,7 @@ use crate::repository::{ get_input_storage_path, persist_input_series_info_batch, MemoryPlaylistSource, PlaylistSource, }; use crate::repository::{xtream_get_file_path, BPlusTreeQuery}; -use crate::utils::ffmpeg::{ProbeFailureKind, ProbeUrlOutcome}; +use crate::utils::ffmpeg::{FfmpegExecutor, ProbeFailureKind, ProbeUrlOutcome}; use crate::utils::{debug_if_enabled, xtream}; use log::{debug, error, info, log_enabled, trace, warn, Level}; use parking_lot::Mutex; @@ -924,7 +924,7 @@ pub async fn update_series_metadata( temp_handle.as_ref().and_then(ProbeHandleGuard::handle), active_handle, ); - match crate::utils::ffmpeg::probe_url_with_cancel( + match FfmpegExecutor::new().probe_url_with_cancel( &episode_url, user_agent.as_deref(), probe_settings.analyze_duration_micros, diff --git a/backend/src/processing/processor/xtream_vod.rs b/backend/src/processing/processor/xtream_vod.rs index 89c6f0321..17026c7ae 100644 --- a/backend/src/processing/processor/xtream_vod.rs +++ b/backend/src/processing/processor/xtream_vod.rs @@ -17,7 +17,7 @@ use crate::ptt::ptt_parse_title; use crate::repository::persist_input_vod_info; use crate::repository::persist_input_vod_info_batch; use crate::repository::{xtream_get_file_path, BPlusTreeQuery}; -use crate::utils::ffmpeg::{ProbeFailureKind, ProbeUrlOutcome}; +use crate::utils::ffmpeg::{FfmpegExecutor, ProbeFailureKind, ProbeUrlOutcome}; use crate::utils::{debug_if_enabled, trace_if_enabled, xtream}; use log::{debug, error, info, log_enabled, trace, warn, Level}; use parking_lot::Mutex; @@ -841,7 +841,7 @@ pub async fn update_vod_metadata( temp_handle.as_ref().and_then(ProbeHandleGuard::handle), active_handle, ); - match crate::utils::ffmpeg::probe_url_with_cancel( + match FfmpegExecutor::new().probe_url_with_cancel( &stream_url, user_agent.as_deref(), analyze_duration, diff --git a/backend/src/repository/user_repository.rs b/backend/src/repository/user_repository.rs index ba9d06415..738f1124b 100644 --- a/backend/src/repository/user_repository.rs +++ b/backend/src/repository/user_repository.rs @@ -368,6 +368,7 @@ pub async fn user_get_bouquet_filter(config: &Config, username: &str, category_i #[cfg(test)] mod tests { use super::*; + use crate::model::MediaToolCapabilities; use crate::utils::FileLockManager; use arc_swap::{ArcSwap, ArcSwapAny}; use shared::model::{ConfigPaths, ProxyType, ProxyUserStatus}; @@ -473,7 +474,7 @@ mod tests { custom_stream_response: Arc::new(ArcSwapAny::default()), access_token_secret: Default::default(), encrypt_secret: Default::default(), - ffprobe_available: Arc::default(), + media_tools: Arc::new(MediaToolCapabilities::new()), }; let target_user = vec![user]; let _ = store_api_user(&cfg, &target_user).await; diff --git a/backend/src/tools/mod.rs b/backend/src/tools/mod.rs deleted file mode 100644 index 6c350a285..000000000 --- a/backend/src/tools/mod.rs +++ /dev/null @@ -1 +0,0 @@ -pub mod lru_cache; diff --git a/backend/src/utils/ffmpeg.rs b/backend/src/utils/ffmpeg.rs index 8bdeed5b0..f3ee92f72 100644 --- a/backend/src/utils/ffmpeg.rs +++ b/backend/src/utils/ffmpeg.rs @@ -2,11 +2,14 @@ use crate::model::ProxyConfig; use log::{debug, warn}; use serde_json::Value; use shared::model::MediaQuality; -use shared::utils::sanitize_sensitive_info; +use shared::utils::{default_thumbnail_height, default_thumbnail_width, sanitize_sensitive_info}; +use std::path::Path; use std::time::Duration; use tokio::process::Command; use url::Url; +const FFMPEG_TIMEOUT: Duration = Duration::from_secs(60); + #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub enum ProbeFailureKind { NotFound, @@ -19,12 +22,220 @@ pub enum ProbeUrlOutcome { Failed(ProbeFailureKind), } -// Checks if ffprobe is available in the system path -pub async fn check_ffprobe_availability() -> bool { - match Command::new("ffprobe").arg("-version").output().await { - Ok(output) => output.status.success(), - Err(_) => false, +#[derive(Debug, Clone, Copy, Default)] +pub struct FfmpegExecutor; + +impl FfmpegExecutor { + #[must_use] + pub const fn new() -> Self { Self } + + /// Checks if the system `ffmpeg` binary is available. + pub async fn check_ffmpeg_availability(&self) -> bool { + self.check_binary_availability("ffmpeg").await } + + // Checks if ffprobe is available in the system path + pub async fn check_ffprobe_availability(&self) -> bool { + self.check_binary_availability("ffprobe").await + } + + /// Extracts a JPEG thumbnail from a local file. + /// Attempts a frame at 180s first and falls back to 0s for short videos. + pub async fn create_thumbnail(&self, input_path: &str, width: u32, height: u32) -> Result, String> { + let temp_dir = tempfile::tempdir() + .map_err(|e| format!("Failed to create temp dir: {e}"))?; + let output_path = temp_dir.path().join("thumb.jpg"); + let scale_filter = build_thumbnail_scale_filter(width, height); + + let output = self.run_ffmpeg_with_timeout(&build_thumbnail_args(input_path, &output_path, &scale_filter, 180)) + .await + .map_err(|e| format!("Failed to run ffmpeg: {e}"))?; + + if !output.status.success() { + let stderr = String::from_utf8_lossy(&output.stderr); + if stderr.contains("Output file is empty") || stderr.contains("nothing was encoded") { + debug!("Video shorter than 180s, retrying at position 0: {input_path}"); + let retry = self.run_ffmpeg_with_timeout(&build_thumbnail_args(input_path, &output_path, &scale_filter, 0)) + .await + .map_err(|e| format!("Failed to run ffmpeg retry: {e}"))?; + + if !retry.status.success() { + let retry_stderr = String::from_utf8_lossy(&retry.stderr); + return Err(format!("ffmpeg failed at position 0: {retry_stderr}")); + } + } else { + return Err(format!("ffmpeg failed: {stderr}")); + } + } + + tokio::fs::read(&output_path) + .await + .map_err(|e| format!("Failed to read thumbnail: {e}")) + } + + pub async fn probe_url( + &self, + url: &str, + user_agent: Option<&str>, + analyze_duration: u64, + probe_size: u64, + timeout_secs: u64, + proxy_cfg: Option<&ProxyConfig>, + ) -> ProbeUrlOutcome { + // Determine timeout: Ensure it's at least as long as the analyze duration + buffer, + // but respect the user setting if it's longer. + let analyze_overhead = Duration::from_micros(analyze_duration) + Duration::from_secs(5); + let config_timeout = Duration::from_secs(timeout_secs); + let timeout_val = std::cmp::max(analyze_overhead, config_timeout); + + let mut command = Command::new("ffprobe"); + + // Ensure the child process is killed if this future is dropped (e.g. by connection preemption) + command.kill_on_drop(true); + + command + .arg("-v").arg("error") + .arg("-show_streams") + .arg("-of").arg("json") + .arg("-analyzeduration").arg(analyze_duration.to_string()) + .arg("-probesize").arg(probe_size.to_string()); + + apply_proxy_to_ffprobe(&mut command, proxy_cfg); + + if let Some(ua) = user_agent { + command.arg("-user_agent").arg(ua); + } + + command.arg(url); + + let output_result = tokio::time::timeout(timeout_val, command.output()).await; + + match output_result { + Ok(Ok(output)) => { + if !output.status.success() { + let stderr = String::from_utf8_lossy(&output.stderr); + debug!("ffprobe failed for {}: {}", sanitize_sensitive_info(url), sanitize_sensitive_info(&stderr)); + if is_not_found_probe_error(&stderr) { + return ProbeUrlOutcome::Failed(ProbeFailureKind::NotFound); + } + return ProbeUrlOutcome::Failed(ProbeFailureKind::Other); + } + + if let Ok(json) = serde_json::from_slice::(&output.stdout) { + if let Some(stream_list) = json.get("streams").and_then(Value::as_array) { + let mut video_stream: Option<&Value> = None; + let mut audio_stream: Option<&Value> = None; + + for stream in stream_list { + let codec_type = stream.get("codec_type").and_then(Value::as_str); + if video_stream.is_none() + && (codec_type == Some("video") + || (codec_type.is_none() + && (stream.get("width").is_some() || stream.get("height").is_some()))) + && !is_attached_pic(stream) + { + video_stream = Some(stream); + } else if audio_stream.is_none() + && (codec_type == Some("audio") + || (codec_type.is_none() + && (stream.get("channels").is_some() + || stream.get("channel_layout").is_some()))) + { + audio_stream = Some(stream); + } + if video_stream.is_some() && audio_stream.is_some() { + break; + } + } + + if video_stream.is_some() || audio_stream.is_some() { + let video_str = video_stream.map(Value::to_string); + let audio_str = audio_stream.map(Value::to_string); + let mq = MediaQuality::from_ffprobe_info(audio_str.as_deref(), video_str.as_deref()); + if let Some(quality) = mq { + return ProbeUrlOutcome::Success( + quality, + video_stream.cloned(), + audio_stream.cloned(), + ); + } + } + } + } else { + warn!("Failed to parse ffprobe json output for {}", sanitize_sensitive_info(url)); + } + } + Ok(Err(e)) => { + warn!("ffprobe execution failed for {}: {}", sanitize_sensitive_info(url), e); + } + Err(_) => { + warn!("ffprobe timed out after {:?} for {}", timeout_val, sanitize_sensitive_info(url)); + } + } + + ProbeUrlOutcome::Failed(ProbeFailureKind::Other) + } + + /// Wrapper around [`Self::probe_url`] that races the probe against an optional cancellation token. + #[allow(clippy::too_many_arguments)] + pub async fn probe_url_with_cancel( + &self, + url: &str, + user_agent: Option<&str>, + analyze_duration: u64, + probe_size: u64, + timeout_secs: u64, + proxy_cfg: Option<&ProxyConfig>, + cancel_token: Option<&tokio_util::sync::CancellationToken>, + ) -> ProbeUrlOutcome { + if let Some(token) = cancel_token { + tokio::select! { + biased; + () = token.cancelled() => { + warn!("Probe preempted for {}", shared::utils::sanitize_sensitive_info(url)); + ProbeUrlOutcome::Failed(ProbeFailureKind::Cancelled) + } + result = self.probe_url(url, user_agent, analyze_duration, probe_size, timeout_secs, proxy_cfg) => result, + } + } else { + self.probe_url(url, user_agent, analyze_duration, probe_size, timeout_secs, proxy_cfg).await + } + } + + async fn check_binary_availability(&self, binary: &str) -> bool { + let mut command = Command::new(binary); + command + .arg("-version") + .kill_on_drop(true); + + match tokio::time::timeout(FFMPEG_TIMEOUT, command.output()).await { + Ok(Ok(output)) => output.status.success(), + Ok(Err(_)) | Err(_) => false, + } + } + + async fn run_ffmpeg_with_timeout(&self, args: &[String]) -> Result { + let child = Command::new("ffmpeg") + .args(args) + .stdout(std::process::Stdio::piped()) + .stderr(std::process::Stdio::piped()) + .kill_on_drop(true) + .spawn() + .map_err(|e| format!("Failed to spawn ffmpeg: {e}"))?; + + tokio::time::timeout(FFMPEG_TIMEOUT, child.wait_with_output()) + .await + .map_err(|_| format_ffmpeg_timeout_error(args))? + .map_err(|e| e.to_string()) + } +} + +fn format_ffmpeg_timeout_error(args: &[String]) -> String { + let summary = args.join(" "); + format!( + "Timed out running ffmpeg after {}s: {summary}", + FFMPEG_TIMEOUT.as_secs() + ) } fn build_ffprobe_proxy_url(proxy_cfg: &ProxyConfig) -> Option { @@ -79,144 +290,48 @@ fn is_attached_pic(stream: &Value) -> bool { fn is_not_found_probe_error(stderr: &str) -> bool { let normalized = stderr.to_ascii_lowercase(); - normalized.contains("404") || normalized.contains("not found") + [ + "http error 404", + "404 not found", + "http/1.1 404", + "http/2 404", + "server returned 404", + ] + .iter() + .any(|needle| normalized.contains(needle)) } -pub async fn probe_url( - url: &str, - user_agent: Option<&str>, - analyze_duration: u64, - probe_size: u64, - timeout_secs: u64, - proxy_cfg: Option<&ProxyConfig>, -) -> ProbeUrlOutcome { - // Determine timeout: Ensure it's at least as long as the analyze duration + buffer, - // but respect the user setting if it's longer. - let analyze_overhead = Duration::from_micros(analyze_duration) + Duration::from_secs(5); - let config_timeout = Duration::from_secs(timeout_secs); - let timeout_val = std::cmp::max(analyze_overhead, config_timeout); - - let mut command = Command::new("ffprobe"); - - // Ensure the child process is killed if this future is dropped (e.g. by connection preemption) - command.kill_on_drop(true); - - command - .arg("-v").arg("error") - .arg("-show_streams") // Get all streams info - .arg("-of").arg("json") - // Optimization for network streams - .arg("-analyzeduration").arg(analyze_duration.to_string()) - .arg("-probesize").arg(probe_size.to_string()); - - apply_proxy_to_ffprobe(&mut command, proxy_cfg); - - if let Some(ua) = user_agent { - command.arg("-user_agent").arg(ua); - } - - command.arg(url); - - let output_result = tokio::time::timeout(timeout_val, command.output()).await; - - match output_result { - Ok(Ok(output)) => { - if !output.status.success() { - let stderr = String::from_utf8_lossy(&output.stderr); - debug!("ffprobe failed for {}: {}", sanitize_sensitive_info(url), sanitize_sensitive_info(&stderr)); - if is_not_found_probe_error(&stderr) { - return ProbeUrlOutcome::Failed(ProbeFailureKind::NotFound); - } - return ProbeUrlOutcome::Failed(ProbeFailureKind::Other); - } - - if let Ok(json) = serde_json::from_slice::(&output.stdout) { - if let Some(stream_list) = json.get("streams").and_then(Value::as_array) { - // Single-pass stream detection: prefer codec_type, fall back to structural hints. - let mut video_stream: Option<&Value> = None; - let mut audio_stream: Option<&Value> = None; - - for stream in stream_list { - let codec_type = stream.get("codec_type").and_then(Value::as_str); - if video_stream.is_none() - && (codec_type == Some("video") - || (codec_type.is_none() - && (stream.get("width").is_some() || stream.get("height").is_some()))) - && !is_attached_pic(stream) - { - video_stream = Some(stream); - } else if audio_stream.is_none() - && (codec_type == Some("audio") - || (codec_type.is_none() - && (stream.get("channels").is_some() - || stream.get("channel_layout").is_some()))) - { - audio_stream = Some(stream); - } - if video_stream.is_some() && audio_stream.is_some() { - break; - } - } - - if video_stream.is_some() || audio_stream.is_some() { - // Materialize strings only for the selected streams. - let video_str = video_stream.map(Value::to_string); - let audio_str = audio_stream.map(Value::to_string); - let mq = MediaQuality::from_ffprobe_info(audio_str.as_deref(), video_str.as_deref()); - if let Some(quality) = mq { - return ProbeUrlOutcome::Success( - quality, - video_stream.cloned(), - audio_stream.cloned(), - ); - } - } - } - } else { - warn!("Failed to parse ffprobe json output for {}", sanitize_sensitive_info(url)); - } - } - Ok(Err(e)) => { - warn!("ffprobe execution failed for {}: {}", sanitize_sensitive_info(url), e); - } - Err(_) => { - warn!("ffprobe timed out after {:?} for {}", timeout_val, sanitize_sensitive_info(url)); - } - } - - ProbeUrlOutcome::Failed(ProbeFailureKind::Other) +fn build_thumbnail_scale_filter(width: u32, height: u32) -> String { + let w = if width < 1 { default_thumbnail_width() } else { width }; + let h = if height < 1 { default_thumbnail_height() } else { height }; + format!( + "scale={w}:{h}:force_original_aspect_ratio=increase,crop={w}:{h}" + ) } -/// Wrapper around [`probe_url`] that races the probe against an optional cancellation token. -/// When the token fires, the probe future is dropped (`kill_on_drop` kills the ffprobe process) -/// and `ProbeUrlOutcome::Failed(Cancelled)` is returned immediately. -pub async fn probe_url_with_cancel( - url: &str, - user_agent: Option<&str>, - analyze_duration: u64, - probe_size: u64, - timeout_secs: u64, - proxy_cfg: Option<&ProxyConfig>, - cancel_token: Option<&tokio_util::sync::CancellationToken>, -) -> ProbeUrlOutcome { - if let Some(token) = cancel_token { - tokio::select! { - biased; - () = token.cancelled() => { - warn!("Probe preempted for {}", shared::utils::sanitize_sensitive_info(url)); - ProbeUrlOutcome::Failed(ProbeFailureKind::Cancelled) - } - result = probe_url(url, user_agent, analyze_duration, probe_size, timeout_secs, proxy_cfg) => result, - } - } else { - probe_url(url, user_agent, analyze_duration, probe_size, timeout_secs, proxy_cfg).await - } +fn build_thumbnail_args(input_path: &str, output_path: &Path, scale_filter: &str, seek_seconds: u32) -> Vec { + vec![ + "-ss".to_string(), + seek_seconds.to_string(), + "-i".to_string(), + input_path.to_string(), + "-frames:v".to_string(), + "1".to_string(), + "-vf".to_string(), + scale_filter.to_string(), + "-q:v".to_string(), + "1".to_string(), + "-y".to_string(), + output_path.to_string_lossy().into_owned(), + ] } #[cfg(test)] mod tests { - use super::build_ffprobe_proxy_url; + use super::{build_ffprobe_proxy_url, build_thumbnail_args, build_thumbnail_scale_filter, format_ffmpeg_timeout_error, FFMPEG_TIMEOUT}; use crate::model::ProxyConfig; + use shared::utils::{default_thumbnail_height, default_thumbnail_width}; + use std::path::Path; #[test] fn build_ffprobe_proxy_url_injects_credentials() { @@ -239,4 +354,60 @@ mod tests { let resolved = build_ffprobe_proxy_url(&proxy_cfg).expect("proxy url should parse"); assert!(resolved.contains("bob:pass@proxy.local:1080")); } + + #[test] + fn build_thumbnail_scale_filter_formats_dimensions() { + let filter = build_thumbnail_scale_filter(320, 180); + assert_eq!(filter, "scale=320:180:force_original_aspect_ratio=increase,crop=320:180"); + } + + #[test] + fn build_thumbnail_scale_filter_uses_defaults_for_zero_dimensions() { + let filter = build_thumbnail_scale_filter(0, 0); + let expected = build_thumbnail_scale_filter(default_thumbnail_width(), default_thumbnail_height()); + assert_eq!(filter, expected); + } + + #[test] + fn build_thumbnail_args_encodes_expected_ffmpeg_call() { + let args = build_thumbnail_args("/tmp/in.mkv", Path::new("/tmp/thumb.jpg"), "scale=320:180", 180); + assert_eq!( + args, + vec![ + "-ss", + "180", + "-i", + "/tmp/in.mkv", + "-frames:v", + "1", + "-vf", + "scale=320:180", + "-q:v", + "1", + "-y", + "/tmp/thumb.jpg", + ] + .into_iter() + .map(str::to_string) + .collect::>() + ); + } + + #[test] + fn format_ffmpeg_timeout_error_includes_binary_timeout_and_args() { + let msg = format_ffmpeg_timeout_error(&["-ss".to_string(), "180".to_string(), "-i".to_string(), "/tmp/in.mkv".to_string()]); + assert!(msg.contains("ffmpeg")); + assert!(msg.contains(&FFMPEG_TIMEOUT.as_secs().to_string())); + assert!(msg.contains("-ss 180 -i /tmp/in.mkv")); + } + + #[test] + fn is_not_found_probe_error_only_matches_http_404_markers() { + assert!(super::is_not_found_probe_error("HTTP error 404 Not Found")); + assert!(super::is_not_found_probe_error("Server returned 404 Not Found")); + assert!(super::is_not_found_probe_error("HTTP/1.1 404 Not Found")); + assert!(!super::is_not_found_probe_error("host not found")); + assert!(!super::is_not_found_probe_error("file not found")); + assert!(!super::is_not_found_probe_error("protocol handler not found")); + } } diff --git a/backend/src/utils/file/config_reader.rs b/backend/src/utils/file/config_reader.rs index 93d5459a0..7e9dc2d0f 100644 --- a/backend/src/utils/file/config_reader.rs +++ b/backend/src/utils/file/config_reader.rs @@ -1,6 +1,6 @@ use crate::api::model::AppState; use crate::model::Config; -use crate::model::{ApiProxyConfig, AppConfig, SourcesConfig}; +use crate::model::{ApiProxyConfig, AppConfig, MediaToolCapabilities, SourcesConfig}; use crate::repository::{ csv_read_inputs, csv_write_inputs, get_api_user_db_path, is_csv_file, load_api_user, }; @@ -29,7 +29,6 @@ use std::io::{self, Read}; use std::path::{Path, PathBuf}; use std::sync::Arc; use tokio::fs; -use tokio::sync::OnceCell; use shared::concat_string; use crate::utils::request::{is_uri}; use url::Url; @@ -581,7 +580,7 @@ pub async fn read_initial_app_config( custom_stream_response: Arc::new(ArcSwapAny::default()), access_token_secret: generate_default_access_secret(), encrypt_secret: generate_default_encrypt_secret(), - ffprobe_available: Arc::new(OnceCell::new()) + media_tools: Arc::new(MediaToolCapabilities::new()), }; app_config.prepare(include_computed)?; //print_info(&app_config); diff --git a/backend/src/tools/lru_cache.rs b/backend/src/utils/lru_cache.rs similarity index 100% rename from backend/src/tools/lru_cache.rs rename to backend/src/utils/lru_cache.rs diff --git a/backend/src/utils/mod.rs b/backend/src/utils/mod.rs index 33eb54245..26784e095 100644 --- a/backend/src/utils/mod.rs +++ b/backend/src/utils/mod.rs @@ -14,6 +14,7 @@ mod db_viewer; mod epg_parser; mod ordinal; pub mod ffmpeg; +mod lru_cache; pub use self::binary_utils::*; pub use self::logging::*; @@ -24,6 +25,7 @@ pub use self::db_viewer::*; pub use shared::utils::*; pub use self::epg_parser::*; pub use self::ordinal::*; +pub use self::lru_cache::*; #[macro_export] macro_rules! debug_if_enabled { diff --git a/backend/src/utils/network/request.rs b/backend/src/utils/network/request.rs index be159c871..bfb43afa4 100644 --- a/backend/src/utils/network/request.rs +++ b/backend/src/utils/network/request.rs @@ -1561,7 +1561,7 @@ mod tests { strip_sensitive_headers_for_cross_origin_redirect, }; use crate::{ - model::{AppConfig, Config, ConfigProvider, ResourceRetryConfig, ReverseProxyConfig, SourcesConfig}, + model::{AppConfig, Config, ConfigProvider, MediaToolCapabilities, ResourceRetryConfig, ReverseProxyConfig, SourcesConfig}, utils::{FileLockManager, DEFAULT_USER_AGENT} }; use arc_swap::{ArcSwap, ArcSwapOption}; @@ -1607,7 +1607,7 @@ mod tests { custom_stream_response: Arc::new(ArcSwapOption::default()), access_token_secret: [0; 32], encrypt_secret: [0; 16], - ffprobe_available: Arc::default(), + media_tools: Arc::new(MediaToolCapabilities::new()), }) } diff --git a/docs/src/configuration/api-proxy.md b/docs/src/configuration/api-proxy.md index b3b4c289e..c020e3283 100644 --- a/docs/src/configuration/api-proxy.md +++ b/docs/src/configuration/api-proxy.md @@ -119,6 +119,7 @@ in your `config.yml`. Without it, these fields are purely cosmetic! | `proxy` | Enum | No | `redirect` | Defines the proxy mode for this user (see [proxy modes](#proxy-modes-proxy) below). | | `server` | String | No | `default` | Which server block (host/port) is rendered into the playlist for this user. | | `epg_timeshift` | String | No | `None` | Shifts EPG times for users in different time zones. Formats supported: `[-+]hh:mm` or `TimeZone`. Examples: `-2:30` (minus 2h30m), `1:45` (1h45m), `+0:15` (15m), `2` (2h), `:30` (30m), `:3` (3m), `Europe/Paris`, `America/New_York`. Only applies when `epg_url` is configured. | +| `epg_request_timeshift` | String | No | `None` | Shifts EPG times for users in different time zones specifically to adjust catchup requests | | `max_connections` | Int | No | `0` | Hard limit of concurrent streams for *this* user. `0` = Unlimited. **Requires** `user_access_control: true` in `config.yml` to be enforced. | | `status` | Enum | No | `Active` | Possible values: `Active`, `Trial`, `Expired`, `Banned`, `Disabled`, `Pending`. **Requires** `user_access_control: true` in `config.yml` to block non-active streaming. | | `exp_date` | UnixTs | No | `None` | Locks the user out after this Unix timestamp. **Requires** `user_access_control: true` in `config.yml` to be enforced. | @@ -137,7 +138,7 @@ This is the most crucial field governing traffic flow for the user. When to use * *Tradeoff:* **No** connection limits, buffering, bandwidth throttling, or custom fallback videos are applied! * **`reverse`**: Tuliprox downloads the video stream from the provider onto your server and pipes it to the client. * *When to use:* This is required for connection limits, fallback videos, caching, bandwidth throttling, and shared - streams to function. + streams to function. * **Partial Syntax**: You can mix and match! `reverse[live]` forces Live-TV through Tuliprox (allowing shared streams) but redirects VODs (saving bandwidth). `reverse[live,vod]` routes everything except Series episodes through Tuliprox. @@ -173,6 +174,324 @@ preempted/killed if any real user needs the slot.)* --- +### EPG Timeshift Configuration Guide + +#### What is an EPG Timeshift? + +**EPG (Electronic Program Guide)** timeshift allows you to adjust TV program times to match your local time zone. This +is especially useful when: + +* You live in a different time zone than your IPTV provider +* You want to view programs as if they were aired at a different time +* Your EPG data is in one time zone, but you need it displayed in another + +--- + +#### The Two EPG Timeshift Fields + +When configuring users in Tuliprox, you'll find two separate fields for EPG timeshift: + +| Field | Purpose | When to Use | +|-------------------------|-----------------------------------------------------------------|--------------------------------------------------------------------| +| `epg_timeshift` | Shifts EPG times for **XMLTV/EPG requests** (regular EPG files) | Use when you want ALL EPG times globally shifted for this user | +| `epg_request_timeshift` | Adjusts client-provided time ranges for **XTream Catchup API** | Use when you need to shift catchup time requests (start/end times) | + +--- + +#### When to Use Which Field? + +##### Use `epg_timeshift` for + +✅ **XMLTV EPG Requests** + +* When clients request EPG via `/xmltv.php` or `/epg` endpoints +* When serving EPG files to applications +* When you want ALL program times adjusted to your time zone + +**Example:** You're in Paris (UTC+2) and your provider is in UTC. All EPG times should be 2 hours earlier. + +--- + +##### Use `epg_request_timeshift` for + +✅ **XTream Catchup API** + +* When clients request catchup via `/timeshift` or `/streaming/timeshift.php` +* When clients provide their own time ranges (`start`, `end`, `duration`) +* When you need to shift those client-provided times to match your needs + +**Example:** A client requests catchup from 14:00-16:00. You want these times shifted to match your local time zone. + +--- + +##### Key Difference + +| Feature | `epg_timeshift` | `epg_request_timeshift` | +|-----------------------|-----------------------------|---------------------------------------------| +| Affects | All EPG program times | Client-provided catchup time ranges only | +| Used in | XMLTV/EPG endpoints | XTream Catchup API | +| Global or Per-Request | Global (affects entire EPG) | Per-request (adjusts client-provided times) | + +--- + +#### Supported Time Formats + +Both fields support the same time formats: + +| Format | Example | Meaning | +|-------------|------------------------------------|-----------------------------------| +| `[-+]hh:mm` | `-2:30`, `+1:00` | Fixed offset in hours and minutes | +| `hh:mm` | `2:00`, `:30` | Positive offset (implies +) | +| `TimeZone` | `Europe/Paris`, `America/New_York` | IANA timezone name | + +##### Common Examples + +| Format | Value | Description | +|---------------|--------------------|-------------------------| +| Fixed offset | `-2:00` | Minus 2 hours | +| | `+1:30` | Plus 1 hour 30 minutes | +| | `2:00` | Plus 2 hours (positive) | +| | `:30` | Plus 30 minutes | +| | `:3` | Plus 3 minutes | +| Timezone name | `Europe/Paris` | Paris time zone | +| | `America/New_York` | New York time zone | +| | `Asia/Tokyo` | Tokyo time zone | + +--- + +#### Configuration Examples + +##### Example 1: User in Paris, No Additional Catchup Offset + +**Scenario:** You're in Paris (UTC+2) and your IPTV provider uses UTC. You want all EPG programs shown in Paris time. +Catchup times should not be adjusted. + +```yaml +# config.yml +api_proxy: + users: + - username: paris_user + password: *** + epg_timeshift: Europe/Paris # ✅ All EPG in Paris time zone + epg_request_timeshift: None # ✅ No adjustment for catchup +``` + +--- + +##### Example 2: User in New York, -1h30 for Everything + +**Scenario:** You're in New York (UTC-5 in winter, UTC-4 in summer). Your EPG is in UTC. You want to shift everything by +-1h30. + +```yaml +# config.yml +api_proxy: + users: + - username: ny_user + password: *** + epg_timeshift: America/New_York # ✅ EPG in New York time zone + epg_request_timeshift: -1:30 # ✅ Catchup also shifted by -1h30 +``` + +--- + +##### Example 3: User in Berlin, -2h Only for Catchup + +**Scenario:** You're in Berlin (UTC+1). Your EPG is already in Berlin time zone. However, when clients request catchup, +you want to adjust their times by -2 hours. + +```yaml +# config.yml +api_proxy: + users: + - username: berlin_user + password: *** + epg_timeshift: Europe/Berlin # ✅ EPG in Berlin time zone + epg_request_timeshift: -2:00 # ✅ Catchup shifted by -2h +``` + +--- + +##### Example 4: User in London, +3h Global, No Catchup Shift + +**Scenario:** You're in London (UTC+0 or UTC+1 depending on DST). Your EPG is 3 hours behind. You want to catch up on +all EPG times. + +```yaml +# config.yml +api_proxy: + users: + - username: london_user + password: *** + epg_timeshift: +3:00 # ✅ All EPG times +3 hours + epg_request_timeshift: None # ✅ Catchup uses client times as-is +``` + +--- + +##### Example 5: User in Sydney, Timezone for EPG, Catchup Unchanged + +**Scenario:** You're in Sydney (UTC+10/UTC+11). You want EPG in Sydney time zone. Catchup should use exact times as +clients request. + +```yaml +# config.yml +api_proxy: + users: + - username: sydney_user + password: *** + epg_timeshift: Australia/Sydney # ✅ EPG in Sydney time zone + epg_request_timeshift: None # ✅ No catchup adjustment +``` + +--- + +#### Real-World Scenarios + +##### Scenario 1: Watching EPG in Different Time Zone + +**Problem:** Your provider's EPG shows programs in UTC time, but you're in Tokyo (UTC+9). + +**Solution:** Use `epg_timeshift: Australia/Sydney` or `epg_timeshift: +9:00`. + +**Result:** All EPG programs appear in Tokyo local time, making it easy to find what's on TV right now. + +--- + +##### Scenario 2: Requesting Catchup for Specific Times + +**Problem:** A client requests catchup from 14:00-16:00 (2pm-4pm), but you want to adjust this for your time zone. + +**Setup:** `epg_request_timeshift: -1:00` + +**When client requests:** `start=14:00&end=16:00` + +**Tuliprox adjusts to:** `start=13:00&end=15:00` (shifted by -1 hour) + +**Result:** Provider receives catchup request for 1pm-3pm (your local time). + +--- + +##### Scenario 3: Same User Needs Different Shifts for Different Features + +**Problem:** You want EPG in your local time zone, but catchup requests should use a different offset (or no offset at +all). + +**Solution:** Configure different values for each field. + +**Result:** XMLTV requests use one shift, XTream catchup uses another shift. Full flexibility! + +--- + +#### Common Questions (FAQ) + +##### Q: Can I leave both fields empty? + +**A:** Yes! Both `epg_timeshift` and `epg_request_timeshift` are optional (`None`). When empty: + +* EPG/EPG times remain unchanged +* Catchup times are exactly as requested by client +* No time shifting is applied + +--- + +##### Q: What's the difference between `+2:00` and `Europe/Paris`? + +**A:** Both are valid formats, but they work differently: + +* `+2:00`: A fixed +2 hour offset. This is always +2 hours, regardless of Daylight Saving Time (DST). +* `Europe/Paris`: Uses the actual Paris time zone, which automatically handles DST (UTC+1 in winter, UTC+2 in summer). + +**Recommendation:** Use timezone names (`Europe/Paris`) for locations with DST. Use fixed offsets (`+2:00`) when you +want a constant shift. + +--- + +##### Q: Which field should I use if I don't know? + +**A:** Start with `epg_timeshift`. This is the most common field and affects XMLTV/EPG requests, which are used by most +IPTV applications. Only configure `epg_request_timeshift` if you specifically need to adjust catchup time requests. + +--- + +##### Q: Can I use both fields with different values? + +**A:** Absolutely! This is the intended design. For example: + +* `epg_timeshift: Europe/Berlin` (EPG in Berlin time) +* `epg_request_timeshift: -2:00` (Catchup shifted by -2 hours) + +Each field affects only its specific use case, giving you complete control. + +--- + +##### Q: How do negative timeshifts work? + +**A:** Negative values shift times **backwards** in time. + +**Example:** `epg_timeshift: -2:30` + +If EPG shows a program at 20:00, it will appear at 17:30 to your client. + +**Common use:** You're ahead of your provider's time zone and need to go back. + +--- + +##### Q: What is the maximum/minimum timeshift I can set? + +**A:** There is no hard limit in Tuliprox, but practical limits apply: + +* **For fixed offsets:** Typically -12 to +14 hours (covering most global time differences) +* **For timezone names:** Any IANA timezone name (e.g., `Pacific/Honolulu` to `Etc/GMT+14`) + +--- + +##### Q: Will timeshift affect recording/catchup? + +**A:** Only `epg_request_timeshift` affects catchup requests. `epg_timeshift` only affects EPG/EPG program times and +does not modify actual stream or recording times. + +--- + +##### Q: Do I need to restart Tuliprox after changing these settings? + +**A:** No! Configuration changes are detected automatically and applied without restart. The new timeshift values will +be used for the next request. + +--- + +#### Quick Reference + +| I Want To... | Use This Field | Format Example | +|--------------------------------------|-------------------------------|----------------------------------------------------------| +| Shift all EPG times globally | `epg_timeshift` | `Europe/Paris` or `-2:00` | +| Adjust catchup time requests | `epg_request_timeshift` | `America/New_York` or `+1:30` | +| No time shifting | Leave both empty | `None` or omit field | +| Handle DST automatically | Use timezone name | `Europe/london` | +| Constant offset regardless of DST | Use fixed offset | `+1:00` | +| EPG in local time, catchup different | Set different values for both | `epg_timeshift: Berlin` / `epg_request_timeshift: -1:00` | + +--- + +#### Getting Help + +If you're unsure which values to use: + +1. **Check your local time zone** and compare with your IPTV provider's EPG +2. **Calculate the difference** in hours (e.g., "I'm 2 hours ahead of EPG") +3. **Set `epg_timeshift` to that value** if you want EPG adjusted +4. **Set `epg_request_timeshift`** only if you specifically need to adjust catchup requests +5. **Test with your IPTV app** to verify times appear correctly + +For timezone names, see: [IANA Time Zone Database](https://en.wikipedia.org/wiki/List_of_tz_database_time_zones) + +--- + +**Note:** These settings only affect EPG (Electronic Program Guide) times. They do not change when streams are actually +aired - that's controlled by your IPTV provider. + +--- +   ## Additional Information diff --git a/frontend/public/assets/i18n/en.json b/frontend/public/assets/i18n/en.json index 66267b8d4..0e55430b4 100644 --- a/frontend/public/assets/i18n/en.json +++ b/frontend/public/assets/i18n/en.json @@ -1233,7 +1233,8 @@ "XTREAM_BATCH": "xc batch", "XTREAM_CODES": "xtream codes", "X_HEADER": "Remove X-* headers", - "YES": "Yes" + "YES": "Yes", + "FORCE": "Force" }, "MESSAGES": { "CLIPBOARD_NOT_SUPPORTED": "Clipboard not supported.\nYour browser or current context does not allow clipboard access.\nPlease use HTTPS or localhost.", diff --git a/frontend/scss/_theme.scss b/frontend/scss/_theme.scss index aa19e09f3..beb962d5e 100644 --- a/frontend/scss/_theme.scss +++ b/frontend/scss/_theme.scss @@ -67,7 +67,7 @@ } :root { - --primary-color: rgba(96, 165, 250, 0.72); + --primary-color: rgba(96, 165, 250, 0.4); /* Icy Blue */ --secondary-color: rgba(251, 113, 133, 0.62); /* Soft Coral Red */ @@ -98,6 +98,16 @@ --card-color: var(--modest-text-color); --sub-card-background-color: rgba(15, 23, 42, 0.62); + --tertiary-background-color: rgba(168, 85, 247, 0.1); + --tertiary-border-color: rgba(168, 85, 247, 0.4); + --tertiary-text-color: rgb(244, 244, 245); + --tertiary-title-text-color: rgb(216, 180, 254); + + --quaternary-background-color: rgba(34, 197, 94, 0.2); + --quaternary-border-color: rgba(34, 197, 94, 0.4); + --quaternary-text-color: rgb(244, 244, 245); + --quaternary-title-text-color: rgb(110, 231, 183); + --quaternary-hover-background-color: rgb(34, 197, 94, 0.3); --text-button-background-color: rgba(30, 64, 175, 0.28); --text-button-color: #ffffff; @@ -114,12 +124,17 @@ --text-button-secondary-hover-background-color: rgba(220, 70, 70, 0.55); --text-button-secondary-hover-color: #ffffff; --text-button-secondary-border-color: rgba(220, 70, 70, 0.5); + --text-button-tertiary-background-color: var(--tertiary-background-color); + --text-button-tertiary-color: var(--tertiary-title-text-color); + --text-button-tertiary-hover-background-color: rgba(168, 85, 247, 0.2); + --text-button-tertiary-hover-color: rgb(233, 213, 255); + --text-button-tertiary-border-color: var(--tertiary-border-color); - --text-button-active-background-color: rgba(34, 197, 94, 0.58); - --text-button-active-color: #fdfdfd; - --text-button-active-hover-background-color: rgba(22, 163, 74, 0.72); - --text-button-active-hover-color: #ffffff; - --text-button-active-border-color: rgba(74, 222, 128, 0.82); + --text-button-active-background-color: var(--quaternary-hover-background-color); + --text-button-active-color: var(--quaternary-title-text-color); + --text-button-active-hover-background-color: var(--quaternary-background-color); + --text-button-active-hover-color: var(--quaternary-text-color); + --text-button-active-border-color: var(--quaternary-border-color); --text-button-disabled-border-color: #717171; --text-button-disabled-background-color: #494949; diff --git a/frontend/scss/app/components/_radio_button_group.scss b/frontend/scss/app/components/_radio_button_group.scss index 5d83c7538..189be4a9d 100644 --- a/frontend/scss/app/components/_radio_button_group.scss +++ b/frontend/scss/app/components/_radio_button_group.scss @@ -1,16 +1,22 @@ .tp__radio-button-group { display: flex; flex-flow: row nowrap; + button { + margin: 0; + } button:not(:first-child) { - border-top-left-radius: 0 !important; - border-bottom-left-radius: 0 !important; - border-right: 0 !important; + border-top-left-radius: 0; + border-bottom-left-radius: 0; } + button:not(:last-child) { - border-top-right-radius: 0 !important; - border-bottom-right-radius: 0 !important; - border-left: 0 !important; + border-top-right-radius: 0; + border-bottom-right-radius: 0; + } + + button:not(:first-child) { + border-left: none; } } \ No newline at end of file diff --git a/frontend/scss/app/components/_tabset.scss b/frontend/scss/app/components/_tabset.scss index 6f71f71e8..f94678754 100644 --- a/frontend/scss/app/components/_tabset.scss +++ b/frontend/scss/app/components/_tabset.scss @@ -26,12 +26,12 @@ } } - .tp__text-button { + .tp__text-button:not(.active) { color: var(--modest-text-color); } .tp__icon-button, - .tp__text-button { + .tp__text-button:not(.active) { &:hover { color: var(--tab-hover-color); fill: currentColor; @@ -40,8 +40,8 @@ } &--active { - .tp__icon-button, - .tp__text-button { + .tp__icon-button:not(.active), + .tp__text-button:not(.active) { background-color: var(--tab-active-background-color); color: var(--tab-active-color); fill: currentColor; diff --git a/frontend/scss/app/components/_text_button.scss b/frontend/scss/app/components/_text_button.scss index 1de5ada44..fe16c70b5 100644 --- a/frontend/scss/app/components/_text_button.scss +++ b/frontend/scss/app/components/_text_button.scss @@ -33,7 +33,7 @@ } } -.tp__text-button.primary { +.tp__text-button.primary:not(.active) { background-color: var(--text-button-primary-background-color); border-color: var(--text-button-primary-border-color); color: var(--text-button-primary-color); @@ -47,7 +47,7 @@ } } -.tp__text-button.secondary { +.tp__text-button.secondary:not(.active) { background-color: var(--text-button-secondary-background-color); border-color: var(--text-button-secondary-border-color); color: var(--text-button-secondary-color); @@ -61,7 +61,24 @@ } } -.tp__text-button.active { +.tp__text-button.tertiary:not(.active) { + background-color: var(--text-button-tertiary-background-color); + border-color: var(--text-button-tertiary-border-color); + color: var(--text-button-tertiary-color); + fill: var(--text-button-tertiary-color); + &:focus, + &:hover { + background-color: var(--text-button-tertiary-hover-background-color); + color: var(--text-button-tertiary-hover-color); + fill: var(--text-button-tertiary-hover-color); + outline: none; + } +} + +.tp__text-button.active, +.tp__text-button.primary.active, +.tp__text-button.secondary.active, +.tp__text-button.tertiary.active { background-color: var(--text-button-active-background-color); border-color: var(--text-button-active-border-color); color: var(--text-button-active-color); @@ -69,6 +86,7 @@ &:focus, &:hover { background-color: var(--text-button-active-hover-background-color); + border-color: var(--text-button-active-border-color); color: var(--text-button-active-hover-color); fill: var(--text-button-active-hover-color); outline: none; diff --git a/frontend/scss/app/components/dashboard/_stats_view.scss b/frontend/scss/app/components/dashboard/_stats_view.scss index 045ef23b0..33fefc28e 100644 --- a/frontend/scss/app/components/dashboard/_stats_view.scss +++ b/frontend/scss/app/components/dashboard/_stats_view.scss @@ -34,6 +34,20 @@ } } } + + &__system { + background-color: var(--tertiary-background-color); + border-color: var(--tertiary-border-color); + + .tp__status-card__body { + color: var(--tertiary-text-color); + } + .tp__status-card__title { + font-weight: bold; + color: var(--tertiary-title-text-color); + } + + } } .tp__collapse-panel.tp__expanded { diff --git a/frontend/src/app/components/card.rs b/frontend/src/app/components/card.rs index a4843cdba..5e3724ba1 100644 --- a/frontend/src/app/components/card.rs +++ b/frontend/src/app/components/card.rs @@ -4,7 +4,7 @@ use yew::prelude::*; #[derive(Properties, Clone, PartialEq, Debug)] pub struct CardProps { #[prop_or_default] - pub class: String, + pub class: Classes, pub children: Children, } @@ -14,7 +14,7 @@ pub fn Card(props: &CardProps) -> Html { let context = CardContext { custom_class: custom_class.clone() }; html! { context={context}> -
+
{ for props.children.iter() }
> diff --git a/frontend/src/app/components/dashboard/stats_view.rs b/frontend/src/app/components/dashboard/stats_view.rs index 9e2a6c17b..352683df8 100644 --- a/frontend/src/app/components/dashboard/stats_view.rs +++ b/frontend/src/app/components/dashboard/stats_view.rs @@ -30,6 +30,16 @@ pub fn StatsView(props: &StatsViewProps) -> Html { }, ); + let render_system_stats = |cache| { + html! { +
+ + + +
+ } + }; + let render_streams_embedded = || { let cache = status_ctx.status.as_ref().map_or_else( || "n/a".to_string(), @@ -44,11 +54,7 @@ pub fn StatsView(props: &StatsViewProps) -> Html {
})}>
-
- - - -
+ { render_system_stats(cache) }
@@ -128,11 +134,7 @@ pub fn StatsView(props: &StatsViewProps) -> Html {

{ translate.t("LABEL.STATS")}

-
- - - -
+ { render_system_stats(cache) }
diff --git a/frontend/src/app/components/playlist/playlist_update_view.rs b/frontend/src/app/components/playlist/playlist_update_view.rs index 4f7dfa0fc..eefbec9f0 100644 --- a/frontend/src/app/components/playlist/playlist_update_view.rs +++ b/frontend/src/app/components/playlist/playlist_update_view.rs @@ -13,7 +13,9 @@ use yew::{platform::spawn_local, prelude::*}; use yew_hooks::use_list; const LABEL_UPDATE_LOCAL_LIBRARY: &str = "LABEL.UPDATE_LOCAL_LIBRARY"; +const LABEL_FORCE: &str = "LABEL.FORCE"; const ACTION_UPDATE_LIBRARY: &str = "update_library"; +const ACTION_UPDATE_LIBRARY_FORCE: &str = "update_library_force"; #[component] pub fn PlaylistUpdateView() -> Html { @@ -84,8 +86,13 @@ pub fn PlaylistUpdateView() -> Html { let services = services.clone(); let translate = translate.clone(); wasm_bindgen_futures::spawn_local(async move { - if name.as_str() == ACTION_UPDATE_LIBRARY { - match services.config.update_library().await { + let mode = match name.as_str() { + ACTION_UPDATE_LIBRARY => 1, + ACTION_UPDATE_LIBRARY_FORCE => 2, + _ => 0, + }; + if mode > 0 { + match services.config.update_library(mode == 2).await { Ok(_) => services.toastr.success(translate.t("MESSAGES.LIBRARY_UPDATE.SUCCESS")), Err(_err) => services.toastr.error(translate.t("MESSAGES.LIBRARY_UPDATE.FAIL")), } @@ -103,10 +110,15 @@ pub fn PlaylistUpdateView() -> Html {

{ translate.t("LABEL.UPDATE")}

{html_if!(can_write_library && library_enabled, { +
+ +
})}
{ html_if!(can_write_playlist, { diff --git a/frontend/src/services/config_service.rs b/frontend/src/services/config_service.rs index fb422b1c6..32a397e9a 100644 --- a/frontend/src/services/config_service.rs +++ b/frontend/src/services/config_service.rs @@ -433,9 +433,9 @@ impl ConfigService { request_get::<()>(&self.geoip_path, None, None).await } - pub async fn update_library(&self) -> Result, Error> { + pub async fn update_library(&self, force_rescan: bool) -> Result, Error> { let path = concat_path(&self.library_path, "scan"); - let params = LibraryScanRequest { force_rescan: false }; + let params = LibraryScanRequest { force_rescan }; request_post::(&path, params, None, None).await } diff --git a/shared/src/model/xtream.rs b/shared/src/model/xtream.rs index b5f901d2d..d88844a9a 100644 --- a/shared/src/model/xtream.rs +++ b/shared/src/model/xtream.rs @@ -254,6 +254,33 @@ pub struct XtreamSeriesInfo { pub episodes: Option>, } +#[cfg(test)] +mod tests { + use super::XtreamVideoInfo; + use serde_json::json; + + #[test] + fn xtream_video_info_accepts_numeric_tmdb_id() { + let parsed: XtreamVideoInfo = serde_json::from_value(json!({ + "info": { + "tmdb_id": 1285728, + "name": "A Normal Woman" + }, + "movie_data": { + "stream_id": 982447, + "name": "A Normal Woman - 2025", + "added": "1771006438", + "category_id": "727", + "container_extension": "mkv", + "direct_source": "" + } + })) + .expect("numeric tmdb_id should deserialize"); + + assert_eq!(parsed.info.tmdb_id.as_ref(), "1285728"); + } +} + // sometimes episodes are a map with season as key, sometimes an array fn deserialize_episodes<'de, D>(deserializer: D) -> Result>, D::Error> where diff --git a/shared/src/utils/string_interner.rs b/shared/src/utils/string_interner.rs index 79c451740..5883dabd4 100644 --- a/shared/src/utils/string_interner.rs +++ b/shared/src/utils/string_interner.rs @@ -134,30 +134,43 @@ pub fn interner_gc() -> usize { 0 } -/// Convert an `f64` that reached `visit_f64` into a round-trip-safe string -/// using the canonical YAML 1.2 spelling (`.inf`, `-.inf`, `.nan`). +/// Convert an `f64` that reached `visit_f64` into a round-trip-safe string. /// -/// `serde_saphyr` recognises these spellings as ambiguous and **quotes** them -/// when re-serializing, so the value survives a YAML round-trip as a string. +/// Special values are emitted as `"infinity"`, `"-infinity"`, and `"nan"`. +/// `serde_saphyr` re-serializes these ambiguous scalars quoted, so they +/// survive a YAML round-trip as strings. /// -/// This is a safety-net: the primary fix is using `deserialize_string` (which -/// skips float parsing entirely), so `visit_f64` is normally not reached for -/// plain string fields. +/// This is a safety-net for paths that intentionally accept typed numeric +/// scalars and normalize them into strings. #[inline] fn f64_to_str(v: f64) -> String { if v.is_infinite() { if v.is_sign_positive() { - ".inf".to_owned() + "infinity".to_owned() } else { - "-.inf".to_owned() + "-infinity".to_owned() } } else if v.is_nan() { - ".nan".to_owned() + "nan".to_owned() } else { v.to_string() } } +#[inline] +// This intentionally only normalizes the lowercase dot-prefixed spellings that +// serde_saphyr emits for special float scalars. Other YAML 1.1 variants such as +// `.Inf`, `.INF`, or `.NaN` are out of scope here because they are not produced +// by the current parser path. +fn normalize_scalar_string(value: &str) -> &str { + match value { + ".inf" => "infinity", + "-.inf" => "-infinity", + ".nan" => "nan", + _ => value, + } +} + // // Two reusable visitor types live here so that multiple public entry-points // can share them without code duplication: @@ -165,9 +178,9 @@ fn f64_to_str(v: f64) -> String { // ArcStrVisitor -> Arc (null/empty -> "") // OptionArcStrVisitor -> Option> (null/empty -> None) // -// `ArcStrVisitor::visit_some` uses `deserialize_string`, which tells saphyr to -// return the **raw scalar text** without float-parsing. That is what makes -// `name: infinity` survive as the literal string `"infinity"`. +// `ArcStrVisitor::visit_some` uses `deserialize_any(self)`, so numeric/bool +// scalars can flow into the typed `visit_*` methods. The tradeoff is that raw +// numeric notation may be normalized (for example `1e2` becomes `"100"`). // // `OptionArcStrVisitor::visit_some` intentionally diverges and uses // `deserialize_any` so JSON/YAML numeric inputs can flow into `visit_i64`, @@ -182,18 +195,19 @@ impl<'de> Visitor<'de> for ArcStrVisitor { fn expecting(&self, f: &mut fmt::Formatter) -> fmt::Result { f.write_str("a string, number, boolean, or null") } - fn visit_str(self, v: &str) -> Result { Ok(v.intern()) } - fn visit_string(self, v: String) -> Result { Ok(v.intern()) } + fn visit_str(self, v: &str) -> Result { + Ok(normalize_scalar_string(v).intern()) + } + fn visit_string(self, v: String) -> Result { + Ok(normalize_scalar_string(v.as_str()).intern()) + } fn visit_bool(self, v: bool) -> Result { Ok(v.to_string().intern()) } fn visit_i64(self, v: i64) -> Result { Ok(v.to_string().intern()) } fn visit_u64(self, v: u64) -> Result { Ok(v.to_string().intern()) } fn visit_f64(self, v: f64) -> Result { Ok(f64_to_str(v).intern()) } fn visit_unit(self) -> Result { Ok("".intern()) } fn visit_none(self) -> Result { Ok("".intern()) } - fn visit_some>(self, d: D) -> Result { - // `deserialize_string` returns the raw text -> `infinity` stays `infinity`. - d.deserialize_string(self) - } + fn visit_some>(self, d: D) -> Result { d.deserialize_any(self) } } /// Visitor that produces `Option>`, mapping null / empty -> `None`. @@ -206,8 +220,22 @@ impl<'de> Visitor<'de> for OptionArcStrVisitor { f.write_str("a string, number, boolean, null, or empty") } - fn visit_str(self, v: &str) -> Result { Ok(Some(v.intern())) } - fn visit_string(self, v: String) -> Result { Ok(Some(v.intern())) } + fn visit_str(self, v: &str) -> Result { + let normalized = normalize_scalar_string(v); + if normalized.is_empty() { + Ok(None) + } else { + Ok(Some(normalized.intern())) + } + } + fn visit_string(self, v: String) -> Result { + let normalized = normalize_scalar_string(v.as_str()); + if normalized.is_empty() { + Ok(None) + } else { + Ok(Some(normalized.intern())) + } + } fn visit_bool(self, v: bool) -> Result { Ok(Some(v.to_string().intern())) } fn visit_i64(self, v: i64) -> Result { Ok(Some(v.to_string().intern())) } fn visit_u64(self, v: u64) -> Result { Ok(Some(v.to_string().intern())) } @@ -217,10 +245,12 @@ impl<'de> Visitor<'de> for OptionArcStrVisitor { fn visit_some>(self, d: D) -> Result { d.deserialize_any(self) } fn visit_seq>(self, mut seq: A) -> Result { while seq.next_element::()?.is_some() {} + log::debug!("ignored sequence while deserializing string interner, returning None"); Ok(None) } fn visit_map>(self, mut map: A) -> Result { while map.next_entry::()?.is_some() {} + log::debug!("ignored map while deserializing string interner, returning None"); Ok(None) } } @@ -259,12 +289,12 @@ pub mod arc_str_serde { serializer.serialize_str(value) } - /// Deserialize a YAML scalar as an interned `Arc`. + /// Deserialize a scalar as an interned `Arc`. /// - /// Uses `deserialize_string` so saphyr hands us the **raw text** without - /// first running it through float/bool/int parsing. This preserves values - /// like `infinity` as the literal string `"infinity"` instead of silently - /// converting them to `".inf"`. + /// This goes through `deserialize_option(ArcStrVisitor)`, and + /// `ArcStrVisitor::visit_some` uses `deserialize_any`. That allows numeric + /// JSON/YAML scalars to be accepted, but raw numeric notation may be lost + /// during normalization (for example `1e2` becomes `"100"`). pub fn deserialize<'de, D>(deserializer: D) -> Result, D::Error> where D: Deserializer<'de>, @@ -312,7 +342,8 @@ pub mod arc_str_option_null_if_empty_serde { // // Reuses `ArcStrVisitor` / `OptionArcStrVisitor` via `deserialize_option`: // - null / ~ / empty -> visit_none / visit_unit -> "" / None -// - `ArcStrVisitor::visit_some` -> deserialize_string -> raw scalar text +// - `ArcStrVisitor::visit_some` -> deserialize_any -> numbers/bools map into +// typed `visit_*` methods, so raw numeric notation may be normalized // - `OptionArcStrVisitor::visit_some` -> deserialize_any -> numbers/bools map // into their typed `visit_*` methods before being interned as strings @@ -336,12 +367,36 @@ where mod tests { use super::*; + #[derive(Debug, serde::Deserialize)] + struct ArcStrHolder { + #[serde(default, with = "arc_str_serde")] + value: Arc, + } + #[derive(Debug, serde::Deserialize)] struct OptArcStrHolder { #[serde(default, with = "arc_str_option_serde")] value: Option>, } + #[test] + fn arc_str_serde_preserves_yaml_infinity_literal_as_string() { + let parsed: ArcStrHolder = serde_saphyr::from_str("value: infinity\n").unwrap(); + assert_eq!(parsed.value.as_ref(), "infinity"); + } + + #[test] + fn arc_str_serde_preserves_yaml_numeric_like_word_as_string() { + let parsed: ArcStrHolder = serde_saphyr::from_str("value: 01abc\n").unwrap(); + assert_eq!(parsed.value.as_ref(), "01abc"); + } + + #[test] + fn arc_str_serde_accepts_json_integer() { + let parsed: ArcStrHolder = serde_json::from_str(r#"{"value":1285728}"#).unwrap(); + assert_eq!(parsed.value.as_ref(), "1285728"); + } + #[test] fn arc_str_option_serde_accepts_json_integer() { let parsed: OptArcStrHolder = serde_json::from_str(r#"{"value":8169}"#).unwrap(); @@ -353,4 +408,38 @@ mod tests { let parsed: OptArcStrHolder = serde_json::from_str(r#"{"value":"8169"}"#).unwrap(); assert_eq!(parsed.value.as_deref(), Some("8169")); } + + #[test] + fn arc_str_option_serde_maps_empty_string_to_none() { + let parsed: OptArcStrHolder = serde_json::from_str(r#"{"value":""}"#).unwrap(); + assert_eq!(parsed.value, None); + } + + #[test] + fn arc_str_serde_normalizes_json_scientific_notation_numbers() { + let parsed: ArcStrHolder = serde_json::from_str(r#"{"value":1e2}"#).unwrap(); + assert_eq!(parsed.value.as_ref(), "100"); + } + + #[test] + fn arc_str_serde_normalizes_yaml_special_float_scalars() { + let parsed_inf: ArcStrHolder = serde_saphyr::from_str("value: .inf\n").unwrap(); + let parsed_neg_inf: ArcStrHolder = serde_saphyr::from_str("value: -.inf\n").unwrap(); + let parsed_nan: ArcStrHolder = serde_saphyr::from_str("value: .nan\n").unwrap(); + + assert_eq!(parsed_inf.value.as_ref(), "infinity"); + assert_eq!(parsed_neg_inf.value.as_ref(), "-infinity"); + assert_eq!(parsed_nan.value.as_ref(), "nan"); + } + + #[test] + fn arc_str_option_serde_normalizes_yaml_special_float_scalars() { + let parsed_inf: OptArcStrHolder = serde_saphyr::from_str("value: .inf\n").unwrap(); + let parsed_neg_inf: OptArcStrHolder = serde_saphyr::from_str("value: -.inf\n").unwrap(); + let parsed_nan: OptArcStrHolder = serde_saphyr::from_str("value: .nan\n").unwrap(); + + assert_eq!(parsed_inf.value.as_deref(), Some("infinity")); + assert_eq!(parsed_neg_inf.value.as_deref(), Some("-infinity")); + assert_eq!(parsed_nan.value.as_deref(), Some("nan")); + } }