From ee0441877d4cf183827bbdde9df18afffb64a7ec Mon Sep 17 00:00:00 2001 From: euzu Date: Fri, 12 Dec 2025 14:35:00 +0100 Subject: [PATCH] xtream and m3u parsing with streaming --- backend/Cargo.toml | 2 +- backend/src/api/api_utils.rs | 13 +- backend/src/api/endpoints/hls_api.rs | 6 +- backend/src/api/endpoints/v1_api.rs | 5 +- backend/src/api/endpoints/v1_api_config.rs | 6 +- backend/src/api/main_api.rs | 2 + backend/src/api/model/app_state.rs | 15 +- backend/src/api/model/streams/mod.rs | 5 +- .../api/model/streams/persist_pipe_stream.rs | 67 +++++- backend/src/model/xmltv.rs | 23 ++- backend/src/processing/parser/m3u.rs | 49 ++--- backend/src/processing/parser/xtream.rs | 52 +++-- backend/src/processing/processor/playlist.rs | 6 +- .../src/processing/processor/xtream_vod.rs | 2 + backend/src/utils/file/csv_input_reader.rs | 2 +- backend/src/utils/geoip.rs | 1 - backend/src/utils/json_utils.rs | 15 ++ backend/src/utils/network/m3u.rs | 6 +- backend/src/utils/network/request.rs | 193 +++++++++++++++--- backend/src/utils/network/xtream.rs | 8 +- shared/src/utils/serde_utils.rs | 63 +++--- 21 files changed, 388 insertions(+), 153 deletions(-) diff --git a/backend/Cargo.toml b/backend/Cargo.toml index cbba13ba1..ebb758e1b 100644 --- a/backend/Cargo.toml +++ b/backend/Cargo.toml @@ -47,7 +47,7 @@ tokio = { version = "1.48", features = ["rt-multi-thread", "parking_lot", "fs"] #console-subscriber = "0" #tracing = "0.1" #tracing-subscriber = { version = "0.3", features = ["fmt", "env-filter"] } -tokio-util = { version = "0.7", features = ["io"] } +tokio-util = { version = "0.7", features = ["io", "io-util"] } tempfile = "3.23" ruzstd = "0.8" filetime = "0.2" diff --git a/backend/src/api/api_utils.rs b/backend/src/api/api_utils.rs index 88615182d..0d2dd0c8a 100644 --- a/backend/src/api/api_utils.rs +++ b/backend/src/api/api_utils.rs @@ -479,13 +479,7 @@ async fn create_stream_response_details( | ProviderStreamState::GracePeriod(_provider_name, request_url) => { let parsed_url = Url::parse(&request_url); let ((stream, stream_info), reconnect_flag) = if let Ok(url) = parsed_url { - let disabled_headers = app_state - .app_config - .config - .load() - .reverse_proxy - .as_ref() - .and_then(|r| r.disabled_header.clone()); + let disabled_headers = app_state.get_disabled_headers(); let provider_stream_factory_options = ProviderStreamFactoryOptions::new( fingerprint.addr, item_type, @@ -1135,10 +1129,7 @@ async fn fetch_resource_with_retry( .reverse_proxy .as_ref() .map_or_else(ResourceRetryConfig::get_default_retry_values, |rp| rp.resource_retry.get_retry_values()); - let disabled_headers = config - .reverse_proxy - .as_ref() - .and_then(|r| r.disabled_header.clone()); + let disabled_headers = app_state.get_disabled_headers(); for attempt in 0..max_attempts { let client = request::get_client_request( &app_state.http_client.load(), diff --git a/backend/src/api/endpoints/hls_api.rs b/backend/src/api/endpoints/hls_api.rs index 57f0fb53b..1f8a890b8 100644 --- a/backend/src/api/endpoints/hls_api.rs +++ b/backend/src/api/endpoints/hls_api.rs @@ -111,11 +111,7 @@ pub(in crate::api) async fn handle_hls_stream_request( // Don't forward Range on playlist fetch; segments use original headers in provider path let filter_header: HeaderFilter = Some(Box::new(|name: &str| !name.eq_ignore_ascii_case("range"))); let forwarded = get_headers_from_request(req_headers, &filter_header); - let config = app_state.app_config.config.load(); - let disabled_headers = config - .reverse_proxy - .as_ref() - .and_then(|r| r.disabled_header.clone()); + let disabled_headers = app_state.get_disabled_headers(); let headers = request::get_request_headers(None, Some(&forwarded), disabled_headers.as_ref()); let input_source = InputSource::from(input).with_url(request_url); match request::download_text_content( diff --git a/backend/src/api/endpoints/v1_api.rs b/backend/src/api/endpoints/v1_api.rs index 39ba81507..558932a34 100644 --- a/backend/src/api/endpoints/v1_api.rs +++ b/backend/src/api/endpoints/v1_api.rs @@ -89,10 +89,7 @@ async fn geoip_update(axum::extract::State(app_state): axum::extract::State { let reader = Cursor::new(content); diff --git a/backend/src/api/endpoints/v1_api_config.rs b/backend/src/api/endpoints/v1_api_config.rs index 8a4a6c64f..7834e430a 100644 --- a/backend/src/api/endpoints/v1_api_config.rs +++ b/backend/src/api/endpoints/v1_api_config.rs @@ -118,11 +118,7 @@ async fn config_batch_content( // The url is changed at this point, we need the raw url for the batch file if let Some(batch_url) = config_input.t_batch_url.as_ref() { let input_source = InputSource::from(&*config_input).with_url(batch_url.to_owned()); - let config = app_state.app_config.config.load(); - let disabled_headers = config - .reverse_proxy - .as_ref() - .and_then(|r| r.disabled_header.clone()); + let disabled_headers = app_state.get_disabled_headers(); return match download_text_content(&app_state.http_client.load(), disabled_headers.as_ref(), &input_source, None, None).await { Ok((content, _path)) => { // Return CSV with explicit content-type diff --git a/backend/src/api/main_api.rs b/backend/src/api/main_api.rs index 6d2948bbd..14d6bf901 100644 --- a/backend/src/api/main_api.rs +++ b/backend/src/api/main_api.rs @@ -31,6 +31,7 @@ use std::sync::atomic::AtomicI8; use std::sync::Arc; use tokio_util::sync::CancellationToken; use tower_governor::key_extractor::SmartIpKeyExtractor; +use tower_http::services::ServeDir; use crate::api::sys_usage::exec_system_usage; use crate::repository::storage::get_geoip_path; use crate::utils::{exec_file_lock_prune, GeoIp}; @@ -315,6 +316,7 @@ pub async fn start_server( // Web Server let mut router = axum::Router::new() .route("/healthcheck", axum::routing::get(healthcheck)) + .nest_service("/.well-known", ServeDir::new(web_dir_path.join("static/.well-known"))) .merge(ws_api_register( web_auth_enabled, web_ui_path.as_str(), diff --git a/backend/src/api/model/app_state.rs b/backend/src/api/model/app_state.rs index 146ed1c8f..d27b894f6 100644 --- a/backend/src/api/model/app_state.rs +++ b/backend/src/api/model/app_state.rs @@ -2,10 +2,7 @@ use crate::api::config_watch::exec_config_watch; use crate::api::model::{ActiveProviderManager, ConnectionManager, EventManager, PlaylistStorage, PlaylistStorageState, SharedStreamManager}; use crate::api::model::{ActiveUserManager, DownloadQueue}; use crate::api::scheduler::exec_scheduler; -use crate::model::{ - AppConfig, Config, ConfigTarget, HdHomeRunConfig, HdHomeRunDeviceConfig, ProcessTargets, - ScheduleConfig, SourcesConfig, -}; +use crate::model::{AppConfig, Config, ConfigTarget, HdHomeRunConfig, HdHomeRunDeviceConfig, ProcessTargets, ReverseProxyDisabledHeaderConfig, ScheduleConfig, SourcesConfig}; use crate::repository::playlist_repository::load_target_into_memory_cache; use crate::tools::lru_cache::LRUResourceCache; use crate::utils::request::create_client; @@ -444,6 +441,16 @@ impl AppState { pub async fn cache_playlist(&self, target_name: &str, playlist: PlaylistStorage) { self.playlists.cache_playlist(target_name, playlist).await; } + + pub fn get_disabled_headers(&self) -> Option { + self + .app_config + .config + .load() + .reverse_proxy + .as_ref() + .and_then(|r| r.disabled_header.clone()) + } } fn schedules_changed(a: &[ScheduleConfig], b: &[ScheduleConfig]) -> bool { diff --git a/backend/src/api/model/streams/mod.rs b/backend/src/api/model/streams/mod.rs index d7a686da4..42aad1a9c 100644 --- a/backend/src/api/model/streams/mod.rs +++ b/backend/src/api/model/streams/mod.rs @@ -5,17 +5,18 @@ mod custom_video_stream; mod transport_stream_buffer; // mod chunked_buffer; mod provider_stream; -mod persist_pipe_stream; mod provider_stream_factory; mod shared_stream_manager; mod active_client_stream; mod throttled_stream; +pub mod persist_pipe_stream; + pub(in crate) use self::transport_stream_buffer::*; pub(in crate::api) use self::provider_stream::*; -pub(in crate::api) use self::persist_pipe_stream::*; pub(in crate::api) use self::provider_stream_factory::*; pub(in crate::api) use self::shared_stream_manager::*; pub(in crate::api) use self::active_client_stream::*; pub(in crate::api) use self::throttled_stream::*; pub(in crate::api) use self::timed_client_stream::*; pub(in crate::api) use self::custom_video_stream::*; +pub use self::persist_pipe_stream::*; diff --git a/backend/src/api/model/streams/persist_pipe_stream.rs b/backend/src/api/model/streams/persist_pipe_stream.rs index 42642a913..32a8c2081 100644 --- a/backend/src/api/model/streams/persist_pipe_stream.rs +++ b/backend/src/api/model/streams/persist_pipe_stream.rs @@ -1,12 +1,13 @@ -use std::path::Path; -use crate::api::model::StreamError; +use crate::utils::request::DynReader; +use crate::utils::{async_file_writer, IO_BUFFER_SIZE}; use bytes::Bytes; use log::{debug, error}; +use std::path::{Path,}; use std::sync::Arc; -use tokio::io::AsyncWriteExt; -use tokio_stream::{StreamExt}; +use tokio::io::{AsyncReadExt, AsyncWriteExt}; use tokio_stream::wrappers::ReceiverStream; -use crate::utils::IO_BUFFER_SIZE; +use tokio_stream::StreamExt; +use crate::api::model::StreamError; pub fn tee_stream( mut stream: S, @@ -14,8 +15,9 @@ pub fn tee_stream( file_path: &Path, callback: Arc, ) -> ReceiverStream> -where S: tokio_stream::Stream> + Send + Unpin + 'static, - W: tokio::io::AsyncWrite + Send + Unpin + 'static, +where + S: tokio_stream::Stream> + Send + Unpin + 'static, + W: tokio::io::AsyncWrite + Send + Unpin + 'static, { let (tx, rx) = tokio::sync::mpsc::channel::>(32); let resource_path = file_path.to_owned(); @@ -77,3 +79,54 @@ where S: tokio_stream::Stream> + Send + Unpin ReceiverStream::new(rx) } + +pub async fn tee_dyn_reader( + reader: DynReader, + persist_path: &Path, + callback: Option>, +) -> DynReader { + let file = match tokio::fs::File::create(persist_path).await { + Ok(f) => f, + Err(err) => { + error!("Cant open file to write: {}, {err}", persist_path.display()); + return reader; + } + }; + + let (mut tx, rx) = tokio::io::duplex(IO_BUFFER_SIZE); + let mut writer = async_file_writer(file); + let reader_arc = reader; + + tokio::spawn(async move { + let mut total_bytes = 0usize; + let mut buf = [0u8; 8192]; + + let mut reader = reader_arc; + + loop { + let n = match reader.read(&mut buf).await { + Ok(0) | Err(_) => break, + Ok(n) => n, + }; + + total_bytes += n; + + if tx.write_all(&buf[..n]).await.is_err() { + break; + } + + if writer.write_all(&buf[..n]).await.is_err() { + break; + } + } + + let _ = writer.flush().await; + let _ = tx.shutdown().await; + + if let Some(cb) = callback { + cb(total_bytes); + } + }); + + Box::pin(rx) as DynReader +} \ No newline at end of file diff --git a/backend/src/model/xmltv.rs b/backend/src/model/xmltv.rs index a47652c6d..138bceb6b 100644 --- a/backend/src/model/xmltv.rs +++ b/backend/src/model/xmltv.rs @@ -12,6 +12,7 @@ use tokio::io::{AsyncRead, AsyncWrite, AsyncWriteExt}; use url::Url; use shared::utils::sanitize_sensitive_info; use crate::api::model::AppState; +use crate::model::{InputSource}; use crate::utils::async_file_reader; use crate::utils::request::{get_remote_content_as_stream}; @@ -211,11 +212,23 @@ pub async fn parse_xmltv_for_web_ui_from_file(path: &Path) -> Result, url: &str) -> Result { if let Ok(request_url) = Url::parse(url) { let client = app_state.http_client.load(); - match get_remote_content_as_stream( - client.as_ref(), - &request_url, - InputFetchMethod::GET, - None, + let input_source: InputSource = InputSource { + name: String::from("xmltv"), + url: request_url.to_string(), + username: None, + password: None, + method: InputFetchMethod::GET, + headers: HashMap::default(), + }; + + let disabled_headers = app_state.get_disabled_headers(); + + match get_remote_content_as_stream( + &client, + &input_source, + None, + &request_url, + disabled_headers.as_ref(), ).await { Ok((stream, _url)) => { parse_xmltv_for_web_ui(stream).await diff --git a/backend/src/processing/parser/m3u.rs b/backend/src/processing/parser/m3u.rs index 5f3bcf4a4..c910d0efd 100644 --- a/backend/src/processing/parser/m3u.rs +++ b/backend/src/processing/parser/m3u.rs @@ -2,7 +2,8 @@ use crate::model::{Config, ConfigInput}; use shared::model::{PlaylistGroup, PlaylistItem, PlaylistItemHeader, PlaylistItemType, XtreamCluster, DEFAULT_VIDEO_EXTENSIONS}; use shared::utils::extract_id_from_url; use std::borrow::BorrowMut; - +use tokio::io::AsyncBufReadExt; +use crate::utils::request::DynReader; // other implementations like calculating text_distance on all titles took too much time // we keep it now as simple as possible and less memory intensive. @@ -87,9 +88,9 @@ fn skip_digit(it: &mut std::str::Chars) -> Option { } } -fn create_empty_playlistitem_header(input_name: &str, url: &str) -> PlaylistItemHeader { +fn create_empty_playlistitem_header(input_name: &str, url: String) -> PlaylistItemHeader { PlaylistItemHeader { - url: url.to_owned(), + url, category_id: 0, input_name: input_name.to_string(), ..Default::default() @@ -107,7 +108,15 @@ macro_rules! process_header_fields { }; } -fn process_header(input_name: &str, video_suffixes: &[&str], content: &str, url: &str) -> PlaylistItemHeader { +fn process_header(input_name: &str, video_suffixes: &[&str], content: &str, url: String) -> PlaylistItemHeader { + let url_id = extract_id_from_url(&url); + let url_types = if video_suffixes.iter().any(|suffix| url.ends_with(suffix)) { + // TODO find Series based on group or configured names + Some((XtreamCluster::Video, PlaylistItemType::Video)) + } else { + None + }; + let mut plih = create_empty_playlistitem_header(input_name, url); let mut it = content.chars(); let mut stack = String::with_capacity(64); @@ -157,7 +166,7 @@ fn process_header(input_name: &str, video_suffixes: &[&str], content: &str, url: plih.epg_channel_id = None; if let Some(pid) = provider_id { plih.id = pid; - } else if let Some(chanid) = extract_id_from_url(url) { + } else if let Some(chanid) = url_id { plih.id = chanid; } } else { @@ -167,11 +176,9 @@ fn process_header(input_name: &str, video_suffixes: &[&str], content: &str, url: } } } - - if video_suffixes.iter().any(|suffix| url.ends_with(suffix)) { - // TODO find Series based on group or configured names - plih.xtream_cluster = XtreamCluster::Video; - plih.item_type = PlaylistItemType::Video; + if let Some((url_cluster, url_item_type)) = url_types { + plih.xtream_cluster = url_cluster; + plih.item_type = url_item_type; } { @@ -189,10 +196,7 @@ fn process_header(input_name: &str, video_suffixes: &[&str], content: &str, url: plih } -pub fn consume_m3u<'a, I, F: FnMut(PlaylistItem)>(cfg: &Config, input: &ConfigInput, lines: I, mut visit: F) -where - I: Iterator, -{ +pub async fn consume_m3u(cfg: &Config, input: &ConfigInput, lines: DynReader, mut visit: F) { let mut header: Option = None; let mut group: Option = None; let input_name = input.name.as_str(); @@ -203,9 +207,10 @@ where }, None => DEFAULT_VIDEO_EXTENSIONS.to_vec() }; - for line in lines { + let mut lines = tokio::io::BufReader::new(lines).lines(); + while let Ok(Some(line)) = lines.next_line().await { if line.starts_with("#EXTINF") { - header = Some(String::from(line)); + header = Some(line); continue; } if line.starts_with("#EXTGRP") { @@ -233,9 +238,7 @@ where } } -pub fn parse_m3u<'a, I>(cfg: &Config, input: &ConfigInput, lines: I) -> Vec -where - I: Iterator, +pub async fn parse_m3u(cfg: &Config, input: &ConfigInput, lines: DynReader) -> Vec { let mut sort_order: Vec> = vec![]; let mut sort_order_idx: usize = 0; @@ -258,7 +261,7 @@ where } } } - }); + }).await; let mut grp_id = 0; let result: Vec = sort_order.into_iter().filter_map(|channels| { // create a group based on the first playlist item @@ -285,7 +288,7 @@ mod test { let url = "http://hello.de/hello.ts"; let line = r#"#EXTINF:-1 channel-id="abc-seven" tvg-id="abc-seven" tvg-logo="https://abc.nz/.images/seven.png" tvg-chno="7" group-title="Sydney" , Seven"#; - let pli = process_header(input, &video_suffixes, line, url); + let pli = process_header(input, &video_suffixes, line, url.to_string()); assert_eq!(pli.title, "Seven"); assert_eq!(pli.id, "abc-seven"); assert_eq!(pli.logo, "https://abc.nz/.images/seven.png"); @@ -300,7 +303,7 @@ mod test { let url = "http://hello.de/hello.ts"; let line = r#"#EXTINF:-1 channel-id="abc-seven" tvg-id="abc-seven" tvg-logo="https://abc.nz/.images/seven.png" tvg-chno="7" group-title="Sydney", Seven"#; - let pli = process_header(input, &video_suffixes, line, url); + let pli = process_header(input, &video_suffixes, line, url.to_string()); assert_eq!(pli.title, "Seven"); assert_eq!(pli.id, "abc-seven"); assert_eq!(pli.logo, "https://abc.nz/.images/seven.png"); @@ -315,7 +318,7 @@ mod test { let url = "http://hello.de/hello.ts"; let line = r#"#EXTINF:-1 tvg-id="abc-seven" xui-id="provider-123" group-title="Sydney", Seven"#; - let pli = process_header(input, &video_suffixes, line, url); + let pli = process_header(input, &video_suffixes, line, url.to_string()); assert_eq!(pli.title, "Seven"); assert_eq!(pli.id, "provider-123"); // Should use xui-id assert_eq!(pli.epg_channel_id, Some("abc-seven".to_string())); // Should preserve original tvg-id diff --git a/backend/src/processing/parser/xtream.rs b/backend/src/processing/parser/xtream.rs index 9a63497f0..40e8074f9 100644 --- a/backend/src/processing/parser/xtream.rs +++ b/backend/src/processing/parser/xtream.rs @@ -7,23 +7,31 @@ use crate::utils::xtream::{get_xtream_stream_url_base}; use serde_json::Value; use std::collections::HashMap; use std::sync::Arc; +use tokio::task::spawn_blocking; +use crate::utils::request::DynReader; -fn map_to_xtream_category(categories: &Value) -> Result, TuliproxError> { - match serde_json::from_value::>(categories.to_owned()) { - Ok(xtream_categories) => Ok(xtream_categories), - Err(err) => { - create_tuliprox_error_result!(TuliproxErrorKind::Notify, "Failed to process categories {}", &err) +async fn map_to_xtream_category(categories: DynReader) -> Result, TuliproxError> { + spawn_blocking(move || { + let reader = tokio_util::io::SyncIoBridge::new(categories); + match serde_json::from_reader::<_, Vec>(reader) { + Ok(xtream_categories) => Ok(xtream_categories), + Err(err) => { + create_tuliprox_error_result!(TuliproxErrorKind::Notify, "Failed to process categories {}", &err) + } } - } + }).await.map_err(|e| TuliproxError::new(TuliproxErrorKind::Notify, format!("Mapping xtream categories failed: {e}")))? } -fn map_to_xtream_streams(xtream_cluster: XtreamCluster, streams: &Value) -> Result, TuliproxError> { - match serde_json::from_value::>(streams.to_owned()) { +async fn map_to_xtream_streams(xtream_cluster: XtreamCluster, streams: DynReader) -> Result, TuliproxError> { + spawn_blocking(move || { + let reader = tokio_util::io::SyncIoBridge::new(streams); + match serde_json::from_reader::<_, Vec>(reader) { Ok(stream_list) => Ok(stream_list), Err(err) => { - create_tuliprox_error_result!(TuliproxErrorKind::Notify, "Failed to map to xtream streams {:?}: {}", xtream_cluster, &err) + create_tuliprox_error_result!(TuliproxErrorKind::Notify, "Failed to map to xtream streams {xtream_cluster}: {err}", ) } } + }).await.map_err(|e| TuliproxError::new(TuliproxErrorKind::Notify, format!("Mapping xtream streams failed: {e}")))? } fn create_xtream_series_episode_url(url: &str, username: &str, password: &str, episode: &XtreamSeriesInfoEpisode) -> Arc { @@ -119,18 +127,18 @@ pub fn create_xtream_url(xtream_cluster: XtreamCluster, url: &str, username: &st } } -pub fn parse_xtream(input: &ConfigInput, +pub async fn parse_xtream(input: &ConfigInput, xtream_cluster: XtreamCluster, - categories: &Value, - streams: &Value) -> Result>, TuliproxError> { - match map_to_xtream_category(categories) { + categories: DynReader, + streams: DynReader) -> Result>, TuliproxError> { + match map_to_xtream_category(categories).await { Ok(xtream_categories) => { let input_name = input.name.clone(); let url = input.url.as_str(); let username = input.username.as_ref().map_or("", |v| v); let password = input.password.as_ref().map_or("", |v| v); - match map_to_xtream_streams(xtream_cluster, streams) { + match map_to_xtream_streams(xtream_cluster, streams).await { Ok(mut xtream_streams) => { let mut group_map: HashMap = xtream_categories.into_iter().map(|category| @@ -204,7 +212,10 @@ pub fn parse_xtream(input: &ConfigInput, #[cfg(test)] mod tests { use std::fs; + use shared::model::XtreamCluster; use crate::model::XtreamSeriesInfo; + use crate::processing::parser::xtream::map_to_xtream_streams; + use crate::utils::async_file_reader; #[test] fn test_read_json_file_into_struct() { @@ -221,4 +232,17 @@ mod tests { } + #[tokio::test] + async fn test_read_json_stream_into_struct() -> std::io::Result<()> { + let reader = Box::pin(async_file_reader(tokio::fs::File::open("/tmp/vod_streams.json").await?)); + match map_to_xtream_streams(XtreamCluster::Video, reader).await { + Ok(_streams) => { + assert!(true); + }, + Err(err) => { + assert!(false, "Failed to parse json file: {err}"); + } + }; + Ok(()) + } } \ No newline at end of file diff --git a/backend/src/processing/processor/playlist.rs b/backend/src/processing/processor/playlist.rs index 710bbf190..198a23647 100644 --- a/backend/src/processing/processor/playlist.rs +++ b/backend/src/processing/processor/playlist.rs @@ -24,7 +24,7 @@ use crate::processing::processor::trakt::process_trakt_categories_for_target; use crate::processing::processor::xtream_series::playlist_resolve_series; use crate::processing::processor::xtream_vod::playlist_resolve_vod; use crate::repository::playlist_repository::persist_playlist; -use crate::utils::debug_if_enabled; +use crate::utils::{debug_if_enabled, trace_if_enabled}; use crate::utils::StepMeasure; use deunicode::deunicode; use futures::StreamExt; @@ -134,7 +134,7 @@ fn exec_rename(pli: &mut PlaylistItem, rename: Option<&Vec>) { let value = get_field_value(result, r.field); let cap = r.pattern.replace_all(value.as_str(), &r.new_name); if log_enabled!(log::Level::Debug) && *value != cap { - debug_if_enabled!("Renamed {}={} to {}", &r.field, value, cap); + trace_if_enabled!("Renamed {}={value} to {cap}", &r.field); } let value = cap.into_owned(); set_field_value(result, r.field, value); @@ -153,7 +153,7 @@ fn rename_playlist(playlist: &mut [PlaylistGroup], target: &ConfigTarget) -> Opt for r in renames { if matches!(r.field, ItemField::Group) { let cap = r.pattern.replace_all(&grp.title, &r.new_name); - debug_if_enabled!("Renamed group {} to {} for {}", &grp.title, cap, target.name); + trace_if_enabled!("Renamed group {} to {cap} for {}", &grp.title, target.name); grp.title = cap.into_owned(); } } diff --git a/backend/src/processing/processor/xtream_vod.rs b/backend/src/processing/processor/xtream_vod.rs index 48500781e..26810540f 100644 --- a/backend/src/processing/processor/xtream_vod.rs +++ b/backend/src/processing/processor/xtream_vod.rs @@ -68,6 +68,8 @@ pub async fn playlist_resolve_vod(app_config: &AppConfig, client: &reqwest::Clie let (resolve_movies, resolve_delay) = get_resolve_vod_options(target, fpl); if !resolve_movies { return; } + // TODO read existing WAL File and import it to avoid duplicate requests + // we cant write to the indexed-document directly because of the write lock and time-consuming operation. // All readers would be waiting for the lock and the app would be unresponsive. // We collect the content into a wal file and write it once we collected everything. diff --git a/backend/src/utils/file/csv_input_reader.rs b/backend/src/utils/file/csv_input_reader.rs index 17fd6faf6..0378d893b 100644 --- a/backend/src/utils/file/csv_input_reader.rs +++ b/backend/src/utils/file/csv_input_reader.rs @@ -190,7 +190,7 @@ pub fn get_csv_file_path(file_uri: &str) -> Result { mod tests { use crate::utils::file::csv_input_reader::csv_read_inputs_from_reader; use crate::utils::{file_reader, resolve_env_var}; - use std::io::{BufReader, Cursor}; + use std::io::{Cursor}; use shared::model::InputType; const M3U_BATCH: &str = r" diff --git a/backend/src/utils/geoip.rs b/backend/src/utils/geoip.rs index e5cadd7d1..07af9ddf8 100644 --- a/backend/src/utils/geoip.rs +++ b/backend/src/utils/geoip.rs @@ -98,7 +98,6 @@ mod test { use crate::utils::geoip::GeoIp; use std::fs::File; - use std::io::BufReader; use std::path::PathBuf; use crate::utils::file_reader; diff --git a/backend/src/utils/json_utils.rs b/backend/src/utils/json_utils.rs index 43888607a..99adfe0ac 100644 --- a/backend/src/utils/json_utils.rs +++ b/backend/src/utils/json_utils.rs @@ -54,3 +54,18 @@ where buf_writer.flush()?; buf_writer.into_inner()?.sync_all() } + +// pub async fn is_valid_json_file(path: &str) -> std::io::Result { +// if let Ok(file) = tokio::fs::File::open(path).await { +// let reader = async_file_reader(file); +// let stream = serde_json::Deserializer::from_reader(reader).into_iter::(); +// for item in stream { +// if item.is_err() { +// return Ok(false); +// } +// } +// Ok(true) +// } else { +// Ok(false) +// } +// } \ No newline at end of file diff --git a/backend/src/utils/network/m3u.rs b/backend/src/utils/network/m3u.rs index d640d177e..dd12c7f42 100644 --- a/backend/src/utils/network/m3u.rs +++ b/backend/src/utils/network/m3u.rs @@ -14,9 +14,9 @@ pub async fn get_m3u_playlist(client: &reqwest::Client, cfg: &Arc, input } }; let persist_file_path = prepare_file_path(input.persist.as_deref(), working_dir, ""); - match request::get_input_text_content(client, &input_source, working_dir, persist_file_path).await { - Ok(text) => { - (m3u::parse_m3u(cfg, input, text.lines()), vec![]) + match request::get_input_text_content_as_stream(client, &input_source, working_dir, persist_file_path).await { + Ok(reader) => { + (m3u::parse_m3u(cfg, input, reader).await, vec![]) } Err(err) => (vec![], vec![err]) } diff --git a/backend/src/utils/network/request.rs b/backend/src/utils/network/request.rs index 9f8baaa82..6e3e14627 100644 --- a/backend/src/utils/network/request.rs +++ b/backend/src/utils/network/request.rs @@ -1,28 +1,30 @@ -use std::collections::{HashMap, HashSet}; -use std::io::{Error, ErrorKind}; -use std::path::{Path, PathBuf}; -use std::pin::Pin; -use std::time::{Duration, Instant}; use futures::{StreamExt, TryStreamExt}; use log::{debug, error, log_enabled, trace, Level}; use reqwest::header::CONTENT_ENCODING; use reqwest::header::{HeaderMap, HeaderName, HeaderValue}; +use std::collections::{HashMap, HashSet}; +use std::io::{Error, ErrorKind}; +use std::path::{Path, PathBuf}; +use std::pin::Pin; +use std::sync::Arc; +use std::time::{Duration, Instant}; use tokio::fs::File; use tokio::io::{AsyncBufReadExt, AsyncRead, AsyncReadExt, AsyncWriteExt}; use tokio_util::io::StreamReader; use url::Url; -use shared::error::create_tuliprox_error_result; -use shared::error::{str_to_io_error, TuliproxError, TuliproxErrorKind}; -use shared::model::{InputFetchMethod, DEFAULT_USER_AGENT}; use crate::model::{format_elapsed_time, AppConfig, InputSource, ReverseProxyDisabledHeaderConfig}; -use crate::model::{ConfigInput}; -use crate::repository::storage::{get_input_storage_path}; +use crate::model::ConfigInput; +use crate::repository::storage::get_input_storage_path; use crate::repository::storage_const; use crate::utils::compression::compression_utils::{is_deflate, is_gzip}; use crate::utils::{async_file_reader, async_file_writer, debug_if_enabled, IO_BUFFER_SIZE}; -use shared::utils::{filter_request_header, sanitize_sensitive_info, short_hash, ENCODING_DEFLATE, ENCODING_GZIP}; use crate::utils::{get_file_path, persist_file}; +use shared::error::create_tuliprox_error_result; +use shared::error::{str_to_io_error, TuliproxError, TuliproxErrorKind}; +use shared::model::{InputFetchMethod, DEFAULT_USER_AGENT}; +use shared::utils::{filter_request_header, human_readable_byte_size, sanitize_sensitive_info, short_hash, ENCODING_DEFLATE, ENCODING_GZIP}; +use crate::api::model::persist_pipe_stream::tee_dyn_reader; #[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)] pub enum MimeCategory { @@ -138,13 +140,63 @@ pub async fn get_input_text_content(client: &reqwest::Client, input: &InputSourc } } + +pub async fn get_input_text_content_as_stream(client: &reqwest::Client, input: &InputSource, working_dir: &str, persist_filepath: Option) -> Result { + debug_if_enabled!("getting input text content working_dir: {}, url: {}", working_dir, sanitize_sensitive_info(&input.url)); + + if input.url.parse::().is_ok() { + match download_text_content_as_stream(client, None, input, None, persist_filepath).await { + Ok((content, _response_url)) => Ok(content), + Err(e) => { + error!("Failed to download input '{}': {}", &input.name, sanitize_sensitive_info(e.to_string().as_str())); + create_tuliprox_error_result!(TuliproxErrorKind::Notify, "Failed to download") + } + } + } else { + let result = match get_file_path(working_dir, Some(PathBuf::from(&input.url))) { + Some(filepath) => { + if filepath.exists() { + match get_local_file_content_as_stream(&filepath).await { + Ok(content) => { + if persist_filepath.is_some() { + let tee_reader: DynReader = if let Some(path) = persist_filepath { + let tee = tee_dyn_reader(content, &path, Some(Arc::new(|size| { + debug_if_enabled!("Persisted {} bytes", human_readable_byte_size(size as u64)); + }))).await; + Box::pin(tee) + } else { + content + }; + Some(tee_reader) + } else { + Some(content) + } + }, + Err(err) => { + return create_tuliprox_error_result!(TuliproxErrorKind::Notify, "Failed : {}", err); + } + } + } else { + None + } + } + None => None + }; + result.map_or_else(|| { + let msg = format!("cant read input url: {}", sanitize_sensitive_info(&input.url)); + error!("{msg}"); + create_tuliprox_error_result!(TuliproxErrorKind::Notify, "{msg}") + }, Ok) + } +} + pub fn get_client_request - (client: &reqwest::Client, - method: InputFetchMethod, - headers: Option<&HashMap>, - url: &Url, - custom_headers: Option<&HashMap, S>>, - disabled_headers: Option<&ReverseProxyDisabledHeaderConfig>) -> reqwest::RequestBuilder { +(client: &reqwest::Client, + method: InputFetchMethod, + headers: Option<&HashMap>, + url: &Url, + custom_headers: Option<&HashMap, S>>, + disabled_headers: Option<&ReverseProxyDisabledHeaderConfig>) -> reqwest::RequestBuilder { let request = match method { InputFetchMethod::GET => client.get(url.clone()), InputFetchMethod::POST => { @@ -233,6 +285,26 @@ pub async fn get_local_file_content(file_path: &Path) -> Result Result { + // open file + let file = File::open(file_path).await.map_err(|err| { + std::io::Error::new(ErrorKind::NotFound, format!("Failed to open file: {}, {err:?}", file_path.display())) + })?; + + let mut buf_reader = async_file_reader(file); + + // Peek first 2 Bytes, for gzip detection + let buffer = buf_reader.fill_buf().await?; + let is_gzipped = buffer.len() >= 2 && is_gzip(&buffer[0..2]); + + if is_gzipped { + // use Async Gzip Decoder + Ok(Box::pin(async_compression::tokio::bufread::GzipDecoder::new(buf_reader))) + } else { + Ok(Box::pin(buf_reader)) + } +} + // pub fn get_local_file_content_blocking(file_path: &PathBuf) -> Result { // match fs::read(file_path) { // Ok(content) => decode_local_file_bytes(content).await, @@ -283,16 +355,23 @@ async fn get_remote_content_as_file(client: &reqwest::Client, input: &ConfigInpu } } -type DynReader = Pin>; +pub type DynReader = Pin>; #[allow(clippy::implicit_hasher)] pub async fn get_remote_content_as_stream( client: &reqwest::Client, + input: &InputSource, + headers: Option<&HeaderMap>, url: &Url, - method: InputFetchMethod, - headers: Option<&HashMap> + disabled_headers: Option<&ReverseProxyDisabledHeaderConfig>, ) -> Result<(DynReader, String), Error> { - let request = get_client_request(client, method, headers, url, None, None); + let custom_headers = headers.map(|h| { + h.iter().map(|(k, v)| (k.as_str().to_string(), v.as_bytes().to_vec())).collect::>() + }); + let merged = get_request_headers(Some(&input.headers), custom_headers.as_ref(), disabled_headers); + let headers: HashMap = merged.iter().map(|(k, v)| (k.as_str().to_string(), String::from_utf8_lossy(v.as_bytes()).to_string())).collect(); + + let request = get_client_request(client, input.method, Some(&headers), url, None, None); let response = request.send().await.map_err(std::io::Error::other)?; if !response.status().is_success() { @@ -331,12 +410,7 @@ pub async fn get_remote_content_as_stream( async fn get_remote_content(client: &reqwest::Client, input: &InputSource, headers: Option<&HeaderMap>, url: &Url, disabled_headers: Option<&ReverseProxyDisabledHeaderConfig>) -> Result<(String, String), Error> { let start_time = Instant::now(); - let custom_headers = headers.map(|h| { - h.iter().map(|(k, v)| (k.as_str().to_string(), v.as_bytes().to_vec())).collect::>()}); - let merged = get_request_headers(Some(&input.headers), custom_headers.as_ref(), disabled_headers); - let headers: HashMap = merged.iter().map(|(k, v)| (k.as_str().to_string(), String::from_utf8_lossy(v.as_bytes()).to_string())).collect(); - - let (mut stream, response_url) = get_remote_content_as_stream(client, url, input.method, Some(&headers)).await.map_err(|e| str_to_io_error(&format!("Failed to read content: {e}")))?; + let (mut stream, response_url) = get_remote_content_as_stream(client, input, headers, url, disabled_headers).await.map_err(|e| str_to_io_error(&format!("Failed to read content: {e}")))?; let mut content = String::new(); stream.read_to_string(&mut content).await.map_err(|e| str_to_io_error(&format!("Failed to read content: {e}")))?; debug_if_enabled!("Request took: {} {}", format_elapsed_time(start_time.elapsed().as_secs()), sanitize_sensitive_info(url.as_str())); @@ -397,10 +471,50 @@ pub async fn download_text_content( Err(err) => Err(err), } } else { - Err(str_to_io_error(&format!( - "Malformed URL {}", - sanitize_sensitive_info(&input.url) - ))) + Err(str_to_io_error(&format!("Malformed URL {}",sanitize_sensitive_info(&input.url)))) + } +} + +pub async fn download_text_content_as_stream( + client: &reqwest::Client, + disabled_headers: Option<&ReverseProxyDisabledHeaderConfig>, + input: &InputSource, + headers: Option<&HeaderMap>, + persist_filepath: Option, +) -> Result<(DynReader, String), Error> { + if let Ok(url) = input.url.parse::() { + let result = if url.scheme() == "file" { + match url.to_file_path() { + Ok(file_path) => get_local_file_content_as_stream(&file_path).await.map(|c| (c, url.to_string())), + Err(()) => Err(str_to_io_error(&format!( + "Unknown file {}", + sanitize_sensitive_info(&input.url) + ))), + } + } else { + get_remote_content_as_stream(client, input, headers, &url, disabled_headers).await + }; + match result { + Ok((content, response_url)) => { + if persist_filepath.is_some() { + + let tee_reader: DynReader = if let Some(path) = persist_filepath { + let tee = tee_dyn_reader(content, &path, Some(Arc::new(|size| { + debug!("Persisted {size} bytes"); + }))).await; + Box::pin(tee) + } else { + content + }; + Ok((tee_reader, response_url)) + } else { + Ok((content, response_url)) + } + } + Err(err) => Err(err), + } + } else { + Err(str_to_io_error(&format!("Malformed URL {}", sanitize_sensitive_info(&input.url)))) } } @@ -424,6 +538,20 @@ pub async fn get_input_json_content(client: &reqwest::Client, disabled_headers: } } +async fn download_json_content_as_stream(client: &reqwest::Client, disabled_headers: Option<&ReverseProxyDisabledHeaderConfig>, input: &InputSource, persist_filepath: Option) -> Result { + debug_if_enabled!("downloading json content from {}", sanitize_sensitive_info(&input.url)); + match download_text_content_as_stream(client, disabled_headers, input, None, persist_filepath).await { + Ok((reader, _response_url)) => Ok(reader), + Err(err) => Err(err) + } +} + +pub async fn get_input_json_content_as_stream(client: &reqwest::Client, disabled_headers: Option<&ReverseProxyDisabledHeaderConfig>, input: &InputSource, persist_filepath: Option) -> Result { + match download_json_content_as_stream(client, disabled_headers, input, persist_filepath).await { + Ok(stream) => Ok(stream), + Err(e) => create_tuliprox_error_result!(TuliproxErrorKind::Notify, "cant download input {}, url: {} => {}", input.name, sanitize_sensitive_info(&input.url), sanitize_sensitive_info(e.to_string().as_str())) + } +} pub fn create_client(cfg: &AppConfig) -> reqwest::ClientBuilder { let config = cfg.config.load(); @@ -433,7 +561,6 @@ pub fn create_client(cfg: &AppConfig) -> reqwest::ClientBuilder { .pool_max_idle_per_host(10) .danger_accept_invalid_certs(config.accept_insecure_ssl_certificates); - if let Some(proxy_cfg) = config.proxy.as_ref() { match Url::parse(&proxy_cfg.url) { Ok(mut url) => { @@ -451,7 +578,7 @@ pub fn create_client(cfg: &AppConfig) -> reqwest::ClientBuilder { Ok(p) => { client = client.proxy(p); } Err(err) => error!("Failed to create SOCKS proxy {url}: {err}"), } - }, + } "http" | "https" => { match reqwest::Proxy::all(url.as_str()) { Ok(p) => { diff --git a/backend/src/utils/network/xtream.rs b/backend/src/utils/network/xtream.rs index 005345bd0..3c713b27b 100644 --- a/backend/src/utils/network/xtream.rs +++ b/backend/src/utils/network/xtream.rs @@ -235,14 +235,14 @@ pub async fn get_xtream_playlist(cfg: &Arc, client: &reqwest::Client, in let stream_file_path = crate::utils::prepare_file_path(input.persist.as_deref(), working_dir, format!("{stream}_").as_str()); match futures::join!( - request::get_input_json_content(client, None, &input_source_category, category_file_path), - request::get_input_json_content(client, None, &input_source_stream, stream_file_path) + request::get_input_json_content_as_stream(client, None, &input_source_category, category_file_path), + request::get_input_json_content_as_stream(client, None, &input_source_stream, stream_file_path) ) { (Ok(category_content), Ok(stream_content)) => { match xtream::parse_xtream(input, *xtream_cluster, - &category_content, - &stream_content) { + category_content, + stream_content).await { Ok(sub_playlist_parsed) => { if let Some(mut xtream_sub_playlist) = sub_playlist_parsed { playlist_groups.append(&mut xtream_sub_playlist); diff --git a/shared/src/utils/serde_utils.rs b/shared/src/utils/serde_utils.rs index 080f55f56..d795df04a 100644 --- a/shared/src/utils/serde_utils.rs +++ b/shared/src/utils/serde_utils.rs @@ -1,7 +1,7 @@ use crate::error::to_io_error; use chrono::{NaiveDateTime, ParseError, TimeZone, Utc}; use serde::de::DeserializeOwned; -use serde::Deserialize; +use serde::{Deserialize, Deserializer}; use serde_json::Value; use std::io; @@ -55,43 +55,52 @@ where }) } -pub fn deserialize_number_from_string<'de, D, T: DeserializeOwned + std::str::FromStr>( - deserializer: D, -) -> Result, D::Error> + +pub fn deserialize_number_from_string<'de, D, T>(deserializer: D) -> Result, D::Error> where - D: serde::Deserializer<'de>, + D: Deserializer<'de>, + T: DeserializeOwned + std::str::FromStr, { - // we define a local enum type inside of the function - // because it is untagged, serde will deserialize as the first variant - // that it can - #[derive(Deserialize)] - #[serde(untagged)] - enum MaybeNumber { - // if it can be parsed as Option, it will be - Value(Option), - // otherwise try parsing as a string - NumberString(String), - } + let raw: Value = Value::deserialize(deserializer)?; - // deserialize into local enum - let value: MaybeNumber = Deserialize::deserialize(deserializer)?; - match value { - // if parsed as T or None, return that - MaybeNumber::Value(value) => Ok(value), + match raw { + // Null → None + Value::Null => Ok(None), - // (if it is any other string) - MaybeNumber::NumberString(s) => { + // its a number + Value::Number(n) => { + let s = n.to_string(); + match s.parse::() { + Ok(v) => Ok(Some(v)), + Err(_) => Ok(None), // Fehler ignorieren, None zurückgeben + } + } + + // String → extract first number + Value::String(s) => { let s = s.trim(); if s.is_empty() { return Ok(None); } - // parse string to number, if fails return None - if let Ok(num) = s.parse::() { - return Ok(Some(num)); + + // find the number + let digits = s.chars() + .skip_while(|c| !c.is_ascii_digit()) + .take_while(|c| c.is_ascii_digit()) + .collect::(); + + if digits.is_empty() { + return Ok(None); } - serde_json::from_str::(s).map_or_else(|_| Ok(None), |val| Ok(Some(val))) + match digits.parse::() { + Ok(v) => Ok(Some(v)), + Err(_) => Ok(None), + } } + + // invalid -> return None + _ => Ok(None), } }