diff --git a/backend/src/api/endpoints/xmltv_api.rs b/backend/src/api/endpoints/xmltv_api.rs index 6cd1f284c..04f4b7f92 100644 --- a/backend/src/api/endpoints/xmltv_api.rs +++ b/backend/src/api/endpoints/xmltv_api.rs @@ -1,12 +1,3 @@ -use axum::response::IntoResponse; -use chrono::{DateTime, Duration, FixedOffset, NaiveDateTime, Offset, TimeZone, Utc}; -use chrono_tz::Tz; -use log::{error, trace}; -use quick_xml::events::{BytesStart, Event}; -use std::path::{Path, PathBuf}; -use std::sync::Arc; -use tokio::io::AsyncWriteExt; -use tokio_util::io::ReaderStream; use crate::api::api_utils::try_unwrap_body; use crate::api::api_utils::{get_user_target, serve_file}; use crate::api::model::AppState; @@ -17,6 +8,15 @@ use crate::repository::m3u_repository::m3u_get_epg_file_path; use crate::repository::storage::get_target_storage_path; use crate::repository::xtream_repository::{xtream_get_epg_file_path, xtream_get_storage_path}; use crate::utils; +use axum::response::IntoResponse; +use chrono::{DateTime, Duration, FixedOffset, NaiveDateTime, Offset, TimeZone, Utc}; +use chrono_tz::Tz; +use log::{error, trace}; +use quick_xml::events::{BytesStart, Event}; +use std::path::{Path, PathBuf}; +use std::sync::Arc; +use tokio::io::AsyncWriteExt; +use tokio_util::io::ReaderStream; pub fn get_empty_epg_response() -> axum::response::Response { try_unwrap_body!(axum::response::Response::builder() @@ -62,7 +62,7 @@ fn time_correct(original: &str, shift: &Duration) -> String { let shifted_dt = dt + *shift; - format!("{} {}",shifted_dt.format("%Y%m%d%H%M%S"), format_offset(tz_offset_minutes)) + format!("{} {}", shifted_dt.format("%Y%m%d%H%M%S"), format_offset(tz_offset_minutes)) } fn format_offset(offset_minutes: i32) -> String { @@ -84,7 +84,7 @@ fn get_epg_path_for_target_of_type(target_name: &str, epg_path: PathBuf) -> Opti None } -pub (in crate::api) fn get_epg_path_for_target(config: &Config, target: &ConfigTarget) -> Option { +pub(in crate::api) fn get_epg_path_for_target(config: &Config, target: &ConfigTarget) -> Option { // TODO if we have multiple targets, first one serves, this can be problematic when // we use m3u playlist but serve xtream target epg @@ -146,95 +146,109 @@ async fn serve_epg( epg_path: &Path, user: &ProxyUserCredentials, ) -> axum::response::Response { - match tokio::fs::File::open(epg_path).await { - Ok(epg_file) => match parse_timeshift(user.epg_timeshift.as_ref()) { - None => serve_file(epg_path, mime::TEXT_XML).await.into_response(), - Some(duration) => serve_epg_with_timeshift(epg_file, duration), - }, + match tokio::fs::try_exists(epg_path).await { + Ok(exists) => { + if exists { + match parse_timeshift(user.epg_timeshift.as_ref()) { + None => serve_file(epg_path, mime::TEXT_XML).await.into_response(), + Some(duration) => serve_epg_with_timeshift(epg_path, duration).await, + } + } else { + get_empty_epg_response() + } + } Err(_) => get_empty_epg_response(), } } -fn serve_epg_with_timeshift( - epg_file: tokio::fs::File, +async fn serve_epg_with_timeshift( + epg_path: &Path, offset_minutes: i32, ) -> axum::response::Response { - let reader = tokio::io::BufReader::new(epg_file); - let (tx, rx) = tokio::io::duplex(8192); - tokio::spawn(async move { - let encoder = async_compression::tokio::write::GzipEncoder::new(tx); - let mut xml_reader = quick_xml::reader::Reader::from_reader(tokio::io::BufReader::new(reader)); - let mut xml_writer = quick_xml::writer::Writer::new(encoder); - let mut buf = Vec::with_capacity(4096); - let duration = Duration::minutes(i64::from(offset_minutes)); + if epg_path.exists() { + match tokio::fs::File::open(epg_path).await { + Ok(file) => { + let reader = tokio::io::BufReader::new(file); + let (tx, rx) = tokio::io::duplex(8192); + tokio::spawn(async move { + let encoder = async_compression::tokio::write::GzipEncoder::new(tx); + let mut xml_reader = quick_xml::reader::Reader::from_reader(tokio::io::BufReader::new(reader)); + let mut xml_writer = quick_xml::writer::Writer::new(encoder); + let mut buf = Vec::with_capacity(4096); + let duration = Duration::minutes(i64::from(offset_minutes)); - loop { - match xml_reader.read_event_into_async(&mut buf).await { - Ok(Event::Start(ref e)) if e.name().as_ref() == b"programme" => { - // Modify the attributes - let mut elem = BytesStart::new(EPG_TAG_PROGRAMME); - for attr in e.attributes() { - match attr { - Ok(attr) if attr.key.as_ref() == b"start" => { - if let Ok(start_value) = attr.decode_and_unescape_value(xml_reader.decoder()) { - // Modify the start attribute value as needed - elem.push_attribute(("start", time_correct(&start_value, &duration).as_str())); - } else { - // keep original attribute unchanged ? - elem.push_attribute(attr); + loop { + match xml_reader.read_event_into_async(&mut buf).await { + Ok(Event::Start(ref e)) if e.name().as_ref() == b"programme" => { + // Modify the attributes + let mut elem = BytesStart::new(EPG_TAG_PROGRAMME); + for attr in e.attributes() { + match attr { + Ok(attr) if attr.key.as_ref() == b"start" => { + if let Ok(start_value) = attr.decode_and_unescape_value(xml_reader.decoder()) { + // Modify the start attribute value as needed + elem.push_attribute(("start", time_correct(&start_value, &duration).as_str())); + } else { + // keep original attribute unchanged ? + elem.push_attribute(attr); + } + } + Ok(attr) if attr.key.as_ref() == b"stop" => { + if let Ok(stop_value) = attr.decode_and_unescape_value(xml_reader.decoder()) { + // Modify the stop attribute value as needed + elem.push_attribute(("stop", time_correct(&stop_value, &duration).as_str())); + } else { + elem.push_attribute(attr); + } + } + Ok(attr) => { + // Copy any other attributes as they are + elem.push_attribute(attr); + } + Err(e) => { + error!("Error parsing attribute: {e}"); + } + } + } + + // Write the modified start event + if let Err(e) = xml_writer.write_event_async(Event::Start(elem)).await { + error!("Failed to write Start event: {e}"); + break; } } - Ok(attr) if attr.key.as_ref() == b"stop" => { - if let Ok(stop_value) = attr.decode_and_unescape_value(xml_reader.decoder()) { - // Modify the stop attribute value as needed - elem.push_attribute(("stop", time_correct(&stop_value, &duration).as_str())); - } else { - elem.push_attribute(attr); + Ok(Event::Eof) => break, // End of file + Ok(event) => { + // Write any other event as is + if let Err(e) = xml_writer.write_event_async(event).await { + error!("Failed to write event: {e}"); + break; } } - Ok(attr) => { - // Copy any other attributes as they are - elem.push_attribute(attr); - } Err(e) => { - error!("Error parsing attribute: {e}"); + error!("Error: {e}"); + break; } } - } - // Write the modified start event - if let Err(e) = xml_writer.write_event_async(Event::Start(elem)).await { - error!("Failed to write Start event: {e}"); - break; + buf.clear(); } - } - Ok(Event::Eof) => break, // End of file - Ok(event) => { - // Write any other event as is - if let Err(e) = xml_writer.write_event_async(event).await { - error!("Failed to write event: {e}"); - break; - } - } - Err(e) => { - error!("Error: {e}"); - break; - } + let _ = xml_writer.into_inner().shutdown().await; + }); + + let body_stream = ReaderStream::new(rx); + try_unwrap_body!(axum::response::Response::builder() + .header( + axum::http::header::CONTENT_TYPE, + mime::TEXT_XML.to_string() + ) + .header(axum::http::header::CONTENT_ENCODING, "gzip") // Set Content-Encoding header + .body(axum::body::Body::from_stream(body_stream))) } - - buf.clear(); - } - let _ = xml_writer.into_inner().shutdown().await; - }); - - let body_stream = ReaderStream::new(rx); - try_unwrap_body!(axum::response::Response::builder() - .header( - axum::http::header::CONTENT_TYPE, - mime::TEXT_XML.to_string() - ) - .header(axum::http::header::CONTENT_ENCODING, "gzip") // Set Content-Encoding header - .body(axum::body::Body::from_stream(body_stream))) + Err(_) => axum::http::StatusCode::INTERNAL_SERVER_ERROR.into_response(), + }; + } + axum::http::StatusCode::NOT_FOUND.into_response() } /// Handles XMLTV EPG API requests, serving the appropriate EPG file with optional time-shifting based on user configuration. diff --git a/backend/src/processing/processor/playlist.rs b/backend/src/processing/processor/playlist.rs index 5d6c86aae..8db849a2d 100644 --- a/backend/src/processing/processor/playlist.rs +++ b/backend/src/processing/processor/playlist.rs @@ -410,12 +410,15 @@ async fn process_sources(client: Arc, config: &Arc, let http_client = Arc::clone(&client); let playlist_state = playlist_state.cloned(); async_tasks.spawn(async move { + // Hold the per-source lock for the full duration of this update. + let current_update_lock = update_lock; let (input_stats, target_stats, mut res_errors) = process_source(Arc::clone(&http_client), cfg, index, usr_trgts, event_manager, playlist_state.as_ref()).await; shared_errors.lock().await.append(&mut res_errors); if let Some(process_stats) = SourceStats::try_new(input_stats, target_stats) { shared_stats.lock().await.push(process_stats); } + drop(current_update_lock); }); } else { let (input_stats, target_stats, mut res_errors) = @@ -424,8 +427,8 @@ async fn process_sources(client: Arc, config: &Arc, if let Some(process_stats) = SourceStats::try_new(input_stats, target_stats) { shared_stats.lock().await.push(process_stats); } + drop(update_lock); } - drop(update_lock); } while let Some(result) = async_tasks.join_next().await { if let Err(err) = result { diff --git a/backend/src/repository/epg_repository.rs b/backend/src/repository/epg_repository.rs index d94da7518..cc2c9feaa 100644 --- a/backend/src/repository/epg_repository.rs +++ b/backend/src/repository/epg_repository.rs @@ -27,31 +27,6 @@ async fn epg_write_file(target: &ConfigTarget, epg: &Epg, path: &Path) -> Result debug_if_enabled!("Epg for target {} written to {}", target.name, path.to_str().unwrap_or("?")); Ok(()) - // - // - // let mut writer = Writer::new(Cursor::new(vec![])); - // match epg.write_to(&mut writer) { - // Ok(()) => { - // let result = writer.into_inner().into_inner(); - // match File::create(path).await { - // Ok(mut epg_file) => { - // match epg_file.write_all("".as_bytes()).await { - // Ok(()) => {} - // Err(err) => return Err(notify_err!(format!("failed to write epg: {} - {}", path.to_str().unwrap_or("?"), err))), - // } - // match epg_file.write_all(&result).await { - // Ok(()) => { - // debug_if_enabled!("Epg for target {} written to {}", target.name, path.to_str().unwrap_or("?")); - // } - // Err(err) => return Err(notify_err!(format!("failed to write epg: {} - {}", path.to_str().unwrap_or("?"), err))), - // } - // } - // Err(err) => return Err(notify_err!(format!("failed to write epg: {} - {}", path.to_str().unwrap_or("?"), err))), - // } - // } - // Err(err) => return Err(notify_err!(format!("failed to write epg: {} - {}", path.to_str().unwrap_or("?"), err))), - // } - // Ok(()) } pub async fn epg_write(cfg: &Config, target: &ConfigTarget, target_path: &Path, epg: Option<&Epg>, output: &TargetOutput) -> Result<(), TuliproxError> { diff --git a/backend/src/repository/user_repository.rs b/backend/src/repository/user_repository.rs index 31a7b8f9b..25e52fc6f 100644 --- a/backend/src/repository/user_repository.rs +++ b/backend/src/repository/user_repository.rs @@ -1,17 +1,17 @@ -use crate::model::{AppConfig, ProxyUserCredentials, TargetUser}; -use crate::model::{Config}; -use shared::model::{PlaylistBouquetDto, ProxyType, ProxyUserStatus, TargetBouquetDto, TargetType, XtreamCluster}; use crate::model::PlaylistXtreamCategory; +use crate::model::{AppConfig, ProxyUserCredentials, TargetUser}; +use crate::model::Config; use crate::repository::bplustree::BPlusTree; use crate::repository::storage_const; use crate::repository::xtream_repository::xtream_get_playlist_categories; +use crate::utils; use crate::utils::json_write_documents_to_file; use chrono::Local; use log::error; +use shared::model::{PlaylistBouquetDto, ProxyType, ProxyUserStatus, TargetBouquetDto, TargetType, XtreamCluster}; use std::collections::{HashMap, HashSet}; -use std::io::{Error}; +use std::io::Error; use std::path::{Path, PathBuf}; -use crate::utils; use tokio::task; #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] @@ -162,7 +162,10 @@ pub async fn store_api_user(cfg: &AppConfig, target_users: &[TargetUser]) -> Res let path = get_api_user_db_path(cfg); backup_api_user_db_file(cfg, &path).await; let write_lock = cfg.file_locks.write_lock(&path).await; - let result = user_tree.store(&path); + let result = task::spawn_blocking({ + let path = path.clone(); + move || user_tree.store(&path) + }).await.map_err(|err| Error::other(format!("Failed to store user db: {err}")))?; drop(write_lock); result } @@ -262,7 +265,7 @@ async fn save_xtream_user_bouquet_for_target(config: &Config, target_name: &str, } if bouquet_path.exists() { - std::fs::remove_file(bouquet_path)?; + tokio::fs::remove_file(bouquet_path).await?; } Ok(()) } @@ -278,7 +281,7 @@ async fn save_m3u_user_bouquet_for_target(storage_path: &Path, target: TargetTyp json_write_documents_to_file(&bouquet_path, bouquet_categories).await?; } None => if bouquet_path.exists() { - std::fs::remove_file(bouquet_path)?; + tokio::fs::remove_file(bouquet_path).await?; } } @@ -409,11 +412,11 @@ pub async fn user_get_bouquet_filter(config: &Config, username: &str, category_i #[cfg(test)] mod tests { use super::*; + use crate::utils::FileLockManager; + use arc_swap::{ArcSwap, ArcSwapAny}; use shared::model::{ConfigPaths, ProxyType, ProxyUserStatus}; use std::env::temp_dir; use std::sync::Arc; - use arc_swap::{ArcSwap, ArcSwapAny}; - use crate::utils::FileLockManager; #[test] pub fn save_target_user() { diff --git a/backend/src/utils/network/request.rs b/backend/src/utils/network/request.rs index cad62f62e..117f91e9d 100644 --- a/backend/src/utils/network/request.rs +++ b/backend/src/utils/network/request.rs @@ -206,8 +206,8 @@ pub fn get_request_headers(request_header pub async fn get_local_file_content(file_path: &Path) -> Result { // Datei öffnen - let file = File::open(file_path).await.map_err(|_| { - std::io::Error::new(ErrorKind::NotFound, format!("File not found: {}", file_path.display())) + 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 = BufReader::new(file); diff --git a/frontend/public/assets/i18n/en.json b/frontend/public/assets/i18n/en.json index 5caeaa687..69c0e56cc 100644 --- a/frontend/public/assets/i18n/en.json +++ b/frontend/public/assets/i18n/en.json @@ -397,7 +397,7 @@ "INFO": { "RESTART_TO_APPLY_CHANGES": "You need to restart to apply changes.", "SCHEDULE_EXAMPLE": "Schedule example", - "THREADS": "Threads: keep it 0 if you have no problems. Increase it if you have multiple providers. Accessing you provider parallel can cause a ban.", + "PROCESS_PARALLEL": "Process Parallel: If you have multiple providers you can enable it. Accessing you provider parallel can cause a ban.", "WEB_ROOT": "Web-Root: Contains the web-ui files. They are normally placed in the same directory under ./web" }, "HINT": {