diff --git a/config/api-proxy.yml b/config/api-proxy.yml index a6de87955..0a5e51276 100644 --- a/config/api-proxy.yml +++ b/config/api-proxy.yml @@ -12,7 +12,7 @@ server: rtmp_port: timezone: Europe/Paris message: Welcome to m3u-filter - path: m3uflt # optional, pnly neeeded for easier reverse proxy configuration see README.md. + path: m3uflt # optional, only needed for easier reverse proxy configuration see README.md. user: - target: pl1 credentials: diff --git a/frontend/src/app/app.tsx b/frontend/src/app/app.tsx index 6c41136dd..ced1f398b 100755 --- a/frontend/src/app/app.tsx +++ b/frontend/src/app/app.tsx @@ -1,4 +1,4 @@ -import React, {useRef, useState, useCallback, useMemo, useEffect} from 'react'; +import React, {useCallback, useEffect, useMemo, useRef, useState} from 'react'; import './app.scss'; import SourceSelector from "../component/source-selector/source-selector"; import PlaylistViewer, {IPlaylistViewer, SearchRequest} from "../component/playlist-viewer/playlist-viewer"; @@ -40,7 +40,7 @@ export default function App(props: AppProps) { setProgress(true); services.playlist().getPlaylist(req).pipe(first()).subscribe({ next: (pl: PlaylistGroup[]) => { - enqueueSnackbar('Sucessfully downloaded playlist', {variant: 'success'}) + enqueueSnackbar('Successfully downloaded playlist', {variant: 'success'}) setPlaylist(pl); }, error: (err) => { diff --git a/src/api/api_utils.rs b/src/api/api_utils.rs index 6d08cfdfc..0463d6340 100644 --- a/src/api/api_utils.rs +++ b/src/api/api_utils.rs @@ -1,19 +1,19 @@ use crate::api::model::app_state::AppState; use crate::api::model::model_utils::get_stream_response_with_headers; -use crate::api::model::persist_pipe_stream::PersistPipeStream; -use crate::api::model::provider_stream; -use crate::api::model::provider_stream::get_provider_pipe_stream; -use crate::api::model::provider_stream_factory::BufferStreamOptions; +use crate::api::model::streams::persist_pipe_stream::PersistPipeStream; +use crate::api::model::streams::provider_stream; +use crate::api::model::streams::provider_stream::get_provider_pipe_stream; +use crate::api::model::streams::provider_stream_factory::BufferStreamOptions; use crate::api::model::request::UserApiRequest; use crate::api::model::stream_error::StreamError; -use crate::{debug_if_enabled, trace_if_enabled}; +use crate::utils::{debug_if_enabled, trace_if_enabled}; use crate::model::api_proxy::ProxyUserCredentials; use crate::model::config::{ConfigInput, ConfigTarget}; use crate::model::playlist::PlaylistItemType; -use crate::utils::file_utils::create_new_file_for_write; -use crate::utils::lru_cache::LRUResourceCache; -use crate::utils::request_utils; -use crate::utils::request_utils::sanitize_sensitive_info; +use crate::utils::file::file_utils::create_new_file_for_write; +use crate::tools::lru_cache::LRUResourceCache; +use crate::utils::network::request; +use crate::utils::network::request::sanitize_sensitive_info; use actix_files::NamedFile; use actix_web::body::{BodyStream, SizedStream}; use actix_web::http::header::{HeaderValue, CACHE_CONTROL}; @@ -26,8 +26,8 @@ use std::collections::HashMap; use std::path::Path; use std::sync::Arc; use url::Url; -use crate::api::model::active_client_stream::ActiveClientStream; -use crate::api::model::shared_stream_manager::SharedStreamManager; +use crate::api::model::streams::active_client_stream::ActiveClientStream; +use crate::api::model::streams::shared_stream_manager::SharedStreamManager; #[macro_export] macro_rules! try_option_bad_request { @@ -67,6 +67,9 @@ macro_rules! try_result_bad_request { }; } +pub use try_option_bad_request; +pub use try_result_bad_request; + pub async fn serve_file(file_path: &Path, req: &HttpRequest, mime_type: mime::Mime) -> HttpResponse { if file_path.exists() { if let Ok(file) = actix_files::NamedFile::open_async(file_path).await { @@ -80,24 +83,24 @@ pub async fn serve_file(file_path: &Path, req: &HttpRequest, mime_type: mime::Mi HttpResponse::NoContent().finish() } -pub fn get_user_target_by_credentials<'a>(username: &str, password: &str, api_req: &'a UserApiRequest, +pub async fn get_user_target_by_credentials<'a>(username: &str, password: &str, api_req: &'a UserApiRequest, app_state: &'a AppState) -> Option<(ProxyUserCredentials, &'a ConfigTarget)> { if !username.is_empty() && !password.is_empty() { - app_state.config.get_target_for_user(username, password) + app_state.config.get_target_for_user(username, password).await } else { let token = api_req.token.as_str().trim(); if token.is_empty() { None } else { - app_state.config.get_target_for_user_by_token(token) + app_state.config.get_target_for_user_by_token(token).await } } } -pub fn get_user_target<'a>(api_req: &'a UserApiRequest, app_state: &'a AppState) -> Option<(ProxyUserCredentials, &'a ConfigTarget)> { +pub async fn get_user_target<'a>(api_req: &'a UserApiRequest, app_state: &'a AppState) -> Option<(ProxyUserCredentials, &'a ConfigTarget)> { let username = api_req.username.as_str().trim(); let password = api_req.password.as_str().trim(); - get_user_target_by_credentials(username, password, api_req, app_state) + get_user_target_by_credentials(username, password, api_req, app_state).await } fn get_stream_options(app_state: &AppState) -> (bool, bool, usize, bool, bool) { @@ -235,7 +238,7 @@ pub async fn resource_response(app_state: &AppState, resource_url: &str, req: &H } trace_if_enabled!("Try to fetch resource {}", sanitize_sensitive_info(resource_url)); if let Ok(url) = Url::parse(resource_url) { - let client = request_utils::get_client_request(&app_state.http_client, input.map(|i| &i.headers), &url, Some(&req_headers)); + let client = request::get_client_request(&app_state.http_client, input.map(|i| &i.headers), &url, Some(&req_headers)); match client.send().await { Ok(response) => { let status = response.status(); diff --git a/src/api/download_api.rs b/src/api/endpoints/download_api.rs similarity index 97% rename from src/api/download_api.rs rename to src/api/endpoints/download_api.rs index e483671bc..c095bd08a 100644 --- a/src/api/download_api.rs +++ b/src/api/endpoints/download_api.rs @@ -1,7 +1,7 @@ use crate::api::model::app_state::AppState; use crate::api::model::download::{DownloadQueue, FileDownload, FileDownloadRequest}; use crate::model::config::VideoDownloadConfig; -use crate::utils::request_utils; +use crate::utils::network::request; use actix_web::{web, HttpResponse}; use async_std::sync::RwLock; use futures::stream::TryStreamExt; @@ -38,7 +38,7 @@ async fn download_file(active: Arc>>, client: &reqwe Err(err) => return Err(format!("Error while writing to file: {file_path_str} {err}")) } } else { - let megabytes = request_utils::bytes_to_megabytes(downloaded); + let megabytes = request::bytes_to_megabytes(downloaded); info!("Downloaded {}, filesize: {}MB", file_path_str, megabytes); active.write().await.as_mut().unwrap().size = downloaded; return Ok(()); @@ -65,7 +65,7 @@ async fn run_download_queue(download_cfg: &VideoDownloadConfig, download_queue: let next_download = download_queue.as_ref().queue.lock().await.pop_front(); if next_download.is_some() { { *download_queue.as_ref().active.write().await = next_download; } - let headers = request_utils::get_request_headers(Some(&download_cfg.headers), None); + let headers = request::get_request_headers(Some(&download_cfg.headers), None); let dq = Arc::clone(download_queue); match reqwest::Client::builder().default_headers(headers).build() { Ok(client) => { diff --git a/src/api/hls_api.rs b/src/api/endpoints/hls_api.rs similarity index 91% rename from src/api/hls_api.rs rename to src/api/endpoints/hls_api.rs index f3c1bdb6d..c62edfb4e 100644 --- a/src/api/hls_api.rs +++ b/src/api/endpoints/hls_api.rs @@ -8,16 +8,16 @@ use crate::api::model::request::UserApiRequest; use crate::model::api_proxy::ProxyUserCredentials; use crate::model::config::{ConfigInput, TargetType}; use crate::model::playlist::{PlaylistEntry, PlaylistItemType, XtreamCluster}; -use crate::processing::hls_parser::{rewrite_hls, M3U_HLSR_PREFIX}; -use crate::{try_option_bad_request, try_result_bad_request}; +use crate::processing::parser::hls::{rewrite_hls, M3U_HLSR_PREFIX}; +use crate::api::api_utils::{try_option_bad_request, try_result_bad_request}; use crate::repository::{m3u_repository, xtream_repository}; use crate::repository::playlist_repository::HLS_EXT; -use crate::utils::request_utils; -use crate::utils::request_utils::{replace_extension, sanitize_sensitive_info}; +use crate::utils::network::request; +use crate::utils::network::request::{replace_extension, sanitize_sensitive_info}; pub(in crate::api) async fn handle_hls_stream_request(app_state: &Data, user: &ProxyUserCredentials, pli: &dyn PlaylistEntry, input: &ConfigInput, target_type: TargetType) -> HttpResponse { let url = replace_extension(&pli.get_provider_url(), HLS_EXT); - match request_utils::download_text_content(Arc::clone(&app_state.http_client), input, &url, None).await { + match request::download_text_content(Arc::clone(&app_state.http_client), input, &url, None).await { Ok(content) => { let hls_content = rewrite_hls(&content, pli.get_virtual_id(), user, &target_type); HttpResponse::Ok().content_type("application/x-mpegurl").body(hls_content) @@ -38,7 +38,7 @@ async fn hls_api_stream( ) -> HttpResponse { let (_token, username, password, channel, _hash, _chunk) = path.into_inner(); let (_user, target) = try_option_bad_request!( - get_user_target_by_credentials(&username, &password, api_req, app_state), + get_user_target_by_credentials(&username, &password, api_req, app_state).await, false, format!("Could not find any user {username}")); diff --git a/src/api/m3u_api.rs b/src/api/endpoints/m3u_api.rs similarity index 91% rename from src/api/m3u_api.rs rename to src/api/endpoints/m3u_api.rs index 545da41b4..6ec41528e 100644 --- a/src/api/m3u_api.rs +++ b/src/api/endpoints/m3u_api.rs @@ -3,8 +3,10 @@ use bytes::Bytes; use futures::stream; use log::{debug, error}; -use crate::api::api_utils::{get_user_target, get_user_target_by_credentials, resource_response, separate_number_and_remainder, stream_response}; -use crate::api::hls_api::handle_hls_stream_request; +use crate::api::api_utils::{get_user_target, get_user_target_by_credentials, resource_response, + separate_number_and_remainder, stream_response, try_option_bad_request, + try_result_bad_request}; +use crate::api::endpoints::hls_api::handle_hls_stream_request; use crate::api::model::app_state::AppState; use crate::api::model::request::UserApiRequest; use crate::model::api_proxy::ProxyType; @@ -12,15 +14,15 @@ use crate::model::config::TargetType; use crate::model::playlist::{FieldGetAccessor, XtreamCluster}; use crate::repository::m3u_playlist_iterator::{M3U_RESOURCE_PATH, M3U_STREAM_PATH}; use crate::repository::m3u_repository::{m3u_get_item_for_stream_id, m3u_load_rewrite_playlist}; -use crate::utils::request_utils::{replace_extension, sanitize_sensitive_info}; -use crate::{debug_if_enabled, try_option_bad_request, try_result_bad_request}; +use crate::utils::network::request::{replace_extension, sanitize_sensitive_info}; +use crate::utils::{debug_if_enabled}; use crate::repository::playlist_repository::HLS_EXT; async fn m3u_api( api_req: &UserApiRequest, app_state: &AppState, ) -> HttpResponse { - match get_user_target(api_req, app_state) { + match get_user_target(api_req, app_state).await { Some((user, target)) => { match m3u_load_rewrite_playlist(&app_state.config, target, &user).await { Ok(m3u_iter) => { @@ -64,7 +66,8 @@ async fn m3u_api_stream( let (username, password, stream_id) = path.into_inner(); let (action_stream_id, stream_ext) = separate_number_and_remainder(&stream_id); let virtual_id: u32 = try_result_bad_request!(action_stream_id.trim().parse()); - let Some((user, target)) = get_user_target_by_credentials(&username, &password, &api_req, &app_state) else { return HttpResponse::BadRequest().finish() }; + let Some((user, target)) = get_user_target_by_credentials(&username, &password, &api_req, &app_state).await + else { return HttpResponse::BadRequest().finish() }; if !target.has_output(&TargetType::M3u) { return HttpResponse::BadRequest().finish(); @@ -104,7 +107,8 @@ async fn m3u_api_resource( ) -> HttpResponse { let (username, password, stream_id, resource) = path.into_inner(); let Ok(m3u_stream_id) = stream_id.parse::() else { return HttpResponse::BadRequest().finish() }; - let Some((user, target)) = get_user_target_by_credentials(&username, &password, &api_req, &app_state) else { return HttpResponse::BadRequest().finish() }; + let Some((user, target)) = get_user_target_by_credentials(&username, &password, &api_req, &app_state).await + else { return HttpResponse::BadRequest().finish() }; if !target.has_output(&TargetType::M3u) { return HttpResponse::BadRequest().finish(); diff --git a/src/api/endpoints/mod.rs b/src/api/endpoints/mod.rs new file mode 100644 index 000000000..f2824be30 --- /dev/null +++ b/src/api/endpoints/mod.rs @@ -0,0 +1,7 @@ +pub(in crate::api) mod download_api; +pub(in crate::api) mod v1_api; +pub(in crate::api) mod xtream_api; +pub(in crate::api) mod m3u_api; +pub(in crate::api) mod xmltv_api; +pub(in crate::api) mod web_index; +pub(in crate::api) mod hls_api; \ No newline at end of file diff --git a/src/api/v1_api.rs b/src/api/endpoints/v1_api.rs similarity index 93% rename from src/api/v1_api.rs rename to src/api/endpoints/v1_api.rs index 2d88a6494..336252c26 100644 --- a/src/api/v1_api.rs +++ b/src/api/endpoints/v1_api.rs @@ -6,7 +6,7 @@ use actix_web_httpauth::middleware::HttpAuthentication; use log::error; use serde_json::json; -use crate::api::download_api; +use crate::api::endpoints::download_api; use crate::api::model::app_state::AppState; use crate::api::model::config::{ServerConfig, ServerInputConfig, ServerSourceConfig, ServerTargetConfig}; use crate::api::model::request::PlaylistRequest; @@ -14,9 +14,10 @@ use crate::auth::authenticator::validator; use crate::m3u_filter_error::M3uFilterError; use crate::model::api_proxy::{ApiProxyConfig, ApiProxyServerInfo, ProxyUserCredentials, TargetUser}; use crate::model::config::{validate_targets, Config, ConfigDto, ConfigInput, ConfigInputOptions, ConfigSource, ConfigTarget, InputType}; -use crate::processing::playlist_processor; -use crate::utils::request_utils::sanitize_sensitive_info; -use crate::utils::{config_reader, m3u_utils, xtream_utils}; +use crate::processing::processor::playlist; +use crate::utils::network::request::sanitize_sensitive_info; +use crate::utils::config_reader; +use crate::utils::network::{m3u, xtream}; fn intern_save_config_api_proxy(backup_dir: &str, api_proxy: &ApiProxyConfig, file_path: &str) -> Option { match config_reader::save_api_proxy(file_path, backup_dir, api_proxy) { @@ -46,7 +47,7 @@ async fn save_config_api_proxy_user( ) -> HttpResponse { let mut users = req.0; users.iter_mut().flat_map(|t| &mut t.credentials).for_each(ProxyUserCredentials::trim); - if let Some(api_proxy) = app_state.config.t_api_proxy.write().unwrap().as_mut() { + if let Some(api_proxy) = app_state.config.t_api_proxy.write().await.as_mut() { let backup_dir = app_state.config.backup_dir.as_ref().unwrap().as_str(); api_proxy.user = users; if let Some(err) = intern_save_config_api_proxy(backup_dir, api_proxy, app_state.config.t_api_proxy_file_path.as_str()) { @@ -84,7 +85,7 @@ async fn save_config_api_proxy_config( return HttpResponse::BadRequest().json(json!({"error": "Invalid content"})); } } - if let Some(api_proxy) = app_state.config.t_api_proxy.write().unwrap().as_mut() { + if let Some(api_proxy) = app_state.config.t_api_proxy.write().await.as_mut() { api_proxy.server = req_api_proxy; let backup_dir = app_state.config.backup_dir.as_ref().unwrap().as_str(); if let Some(err) = intern_save_config_api_proxy(backup_dir, api_proxy, app_state.config.t_api_proxy_file_path.as_str()) { @@ -103,7 +104,7 @@ async fn playlist_update( let process_targets = validate_targets(user_targets.as_ref(), &app_state.config.sources); match process_targets { Ok(valid_targets) => { - actix_rt::spawn(playlist_processor::exec_processing(Arc::clone(&app_state.http_client), Arc::clone(&app_state.config), Arc::new(valid_targets))); + actix_rt::spawn(playlist::exec_processing(Arc::clone(&app_state.http_client), Arc::clone(&app_state.config), Arc::new(valid_targets))); HttpResponse::Ok().finish() } Err(err) => { @@ -136,8 +137,8 @@ async fn get_playlist(client: Arc, cfg_input: Option<&ConfigInp Some(input) => { let (result, errors) = match input.input_type { - InputType::M3u => m3u_utils::get_m3u_playlist(client, cfg, input, &cfg.working_dir).await, - InputType::Xtream => xtream_utils::get_xtream_playlist(client, input, &cfg.working_dir).await, + InputType::M3u => m3u::get_m3u_playlist(client, cfg, input, &cfg.working_dir).await, + InputType::Xtream => xtream::get_xtream_playlist(client, input, &cfg.working_dir).await, }; if result.is_empty() { let error_strings: Vec = errors.iter().map(std::string::ToString::to_string).collect(); @@ -226,7 +227,7 @@ async fn config( // if we didn't read it from file then we should use it from app_state if result.api_proxy.is_none() { - result.api_proxy.clone_from(&app_state.config.t_api_proxy.read().unwrap()); + result.api_proxy.clone_from(&*app_state.config.t_api_proxy.read().await); } HttpResponse::Ok().json(result) diff --git a/src/api/web_index.rs b/src/api/endpoints/web_index.rs similarity index 100% rename from src/api/web_index.rs rename to src/api/endpoints/web_index.rs diff --git a/src/api/xmltv_api.rs b/src/api/endpoints/xmltv_api.rs similarity index 97% rename from src/api/xmltv_api.rs rename to src/api/endpoints/xmltv_api.rs index cad9bf0e4..ef7bc9264 100644 --- a/src/api/xmltv_api.rs +++ b/src/api/endpoints/xmltv_api.rs @@ -1,7 +1,7 @@ use std::fs::File; use std::path::{Path, PathBuf}; -use actix_web::{HttpRequest, HttpResponse, web, http::header}; +use actix_web::{http::header, web, HttpRequest, HttpResponse}; use log::{error, trace}; use quick_xml::{Reader, Writer}; use flate2::write::GzEncoder; @@ -12,14 +12,14 @@ use chrono::{Duration, NaiveDateTime, TimeDelta}; use crate::api::api_utils::{get_user_target, serve_file}; use crate::api::model::app_state::AppState; use crate::api::model::request::UserApiRequest; -use crate::model::api_proxy::{ProxyUserCredentials}; +use crate::model::api_proxy::ProxyUserCredentials; use crate::model::config::{Config, ConfigTarget}; use crate::model::config::TargetType; 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::{file_utils}; -use crate::utils::file_utils::file_reader; +use crate::utils::file::file_utils; +use crate::utils::file::file_utils::file_reader; fn time_correct(date_time: &str, correction: &TimeDelta) -> String { // Split the dateTime string into date and time parts @@ -163,7 +163,7 @@ async fn xmltv_api( req: HttpRequest, app_state: web::Data, ) -> HttpResponse { - if let Some((user, target)) = get_user_target(&api_req, &app_state) { + if let Some((user, target)) = get_user_target(&api_req, &app_state).await { match get_epg_path_for_target(&app_state.config, target) { None => { // No epg configured, No processing or timeshift, epg can't be mapped to the channels. diff --git a/src/api/xtream_api.rs b/src/api/endpoints/xtream_api.rs similarity index 93% rename from src/api/xtream_api.rs rename to src/api/endpoints/xtream_api.rs index 7d646af05..d1ae2a501 100644 --- a/src/api/xtream_api.rs +++ b/src/api/endpoints/xtream_api.rs @@ -1,12 +1,13 @@ // https://github.com/tellytv/go.xtream-codes/blob/master/structs.go -use crate::{trace_if_enabled, try_option_bad_request, try_result_bad_request, Arc}; +use crate::api::api_utils::{try_option_bad_request, try_result_bad_request}; +use crate::utils::{trace_if_enabled}; use std::collections::HashMap; use std::fmt::{Display, Formatter}; use std::path::Path; use std::rc::Rc; use std::str::FromStr; - +use std::sync::Arc; use actix_web::{web, HttpRequest, HttpResponse}; use bytes::Bytes; use futures::stream::{self, StreamExt}; @@ -15,7 +16,7 @@ use log::{debug, error, warn}; use serde_json::{Map, Value}; use crate::api::api_utils::{get_user_target, get_user_target_by_credentials, resource_response, separate_number_and_remainder, serve_file, stream_response}; -use crate::api::hls_api::handle_hls_stream_request; +use crate::api::endpoints::hls_api::handle_hls_stream_request; use crate::api::model::app_state::AppState; use crate::api::model::request::UserApiRequest; use crate::api::model::xtream::XtreamAuthorizationResponse; @@ -30,12 +31,14 @@ use crate::repository::storage::{get_target_storage_path, hex_encode}; use crate::repository::target_id_mapping::TargetIdMapping; use crate::repository::xtream_repository; use crate::repository::xtream_repository::{TAG_EPISODES, TAG_INFO_DATA, TAG_SEASONS_DATA}; -use crate::utils::hash_utils::{generate_playlist_uuid}; -use crate::utils::json_utils::{get_u32_from_serde_value}; -use crate::utils::request_utils::{extract_extension_from_url, replace_extension, sanitize_sensitive_info}; -use crate::utils::xtream_utils::{create_vod_info_from_item, ACTION_GET_LIVE_CATEGORIES, ACTION_GET_LIVE_STREAMS, ACTION_GET_SERIES, ACTION_GET_SERIES_CATEGORIES, ACTION_GET_SERIES_INFO, ACTION_GET_VOD_CATEGORIES, ACTION_GET_VOD_INFO, ACTION_GET_VOD_STREAMS}; -use crate::utils::{json_utils, request_utils, xtream_utils}; -use crate::{debug_if_enabled, info_err}; +use crate::utils::hash_utils::generate_playlist_uuid; +use crate::utils::json_utils::get_u32_from_serde_value; +use crate::utils::network::request::{extract_extension_from_url, replace_extension, sanitize_sensitive_info}; +use crate::utils::network::xtream::{create_vod_info_from_item, ACTION_GET_LIVE_CATEGORIES, ACTION_GET_LIVE_STREAMS, ACTION_GET_SERIES, ACTION_GET_SERIES_CATEGORIES, ACTION_GET_SERIES_INFO, ACTION_GET_VOD_CATEGORIES, ACTION_GET_VOD_INFO, ACTION_GET_VOD_STREAMS}; +use crate::utils::json_utils; +use crate::utils::debug_if_enabled; +use crate::m3u_filter_error::info_err; +use crate::utils::network::{request, xtream}; const ACTION_GET_EPG: &str = "get_epg"; const ACTION_GET_SHORT_EPG: &str = "get_short_epg"; @@ -129,9 +132,8 @@ fn get_xtream_player_api_stream_url(input: &ConfigInput, context: &XtreamApiStre } } -fn get_user_info(user: &ProxyUserCredentials, cfg: &Config) -> XtreamAuthorizationResponse { - let server_info = cfg.get_user_server_info(user); - +async fn get_user_info(user: &ProxyUserCredentials, cfg: &Config) -> XtreamAuthorizationResponse { + let server_info = cfg.get_user_server_info(user).await; XtreamAuthorizationResponse::new(&server_info, user) } @@ -141,7 +143,7 @@ async fn xtream_player_api_stream( app_state: &web::Data, stream_req: XtreamApiStreamRequest<'_>, ) -> HttpResponse { - let (user, target) = try_option_bad_request!(get_user_target_by_credentials(stream_req.username, stream_req.password, api_req, app_state), false, format!("Could not find any user {}", stream_req.username)); + let (user, target) = try_option_bad_request!(get_user_target_by_credentials(stream_req.username, stream_req.password, api_req, app_state).await, false, format!("Could not find any user {}", stream_req.username)); let target_name = &target.name; if !target.has_output(&TargetType::Xtream) { debug!("Target has no xtream output {}", target_name); @@ -329,7 +331,7 @@ async fn xtream_player_api_resource( app_state: &web::Data, resource_req: XtreamApiStreamRequest<'_>, ) -> HttpResponse { - let (user, target) = try_option_bad_request!(get_user_target_by_credentials(resource_req.username, resource_req.password, api_req, app_state), false, format!("Could not find any user {}", resource_req.username)); + let (user, target) = try_option_bad_request!(get_user_target_by_credentials(resource_req.username, resource_req.password, api_req, app_state).await, false, format!("Could not find any user {}", resource_req.username)); let target_name = &target.name; if !target.has_output(&TargetType::Xtream) { debug!("Target has no xtream output {}", target_name); @@ -436,11 +438,11 @@ async fn xtream_get_stream_info_response(app_state: &AppState, user: &ProxyUserC if pli.provider_id > 0 { let input_name = Rc::clone(&pli.input_name); if let Some(input) = app_state.config.get_input_by_name(input_name.as_str()) { - if let Some(info_url) = xtream_utils::get_xtream_player_api_info_url(input, cluster, pli.provider_id) { + if let Some(info_url) = xtream::get_xtream_player_api_info_url(input, cluster, pli.provider_id) { // Redirect is only possible for live streams, vod and series info needs to be modified if user.proxy == ProxyType::Redirect && cluster == XtreamCluster::Live { return HttpResponse::Found().insert_header(("Location", info_url)).finish(); - } else if let Ok(content) = xtream_utils::get_xtream_stream_info(Arc::clone(&app_state.http_client), &app_state.config, user, input, target, &pli, info_url.as_str(), cluster).await { + } else if let Ok(content) = xtream::get_xtream_stream_info(Arc::clone(&app_state.http_client), &app_state.config, user, input, target, &pli, info_url.as_str(), cluster).await { return HttpResponse::Ok().content_type(mime::APPLICATION_JSON).body(content); } } @@ -470,7 +472,7 @@ async fn xtream_get_short_epg(app_state: &AppState, user: &ProxyUserCredentials, if pli.provider_id > 0 { let input_name = Rc::clone(&pli.input_name); if let Some(input) = app_state.config.get_input_by_name(input_name.as_str()) { - if let Some(action_url) = xtream_utils::get_xtream_player_api_action_url(input, ACTION_GET_SHORT_EPG) { + if let Some(action_url) = xtream::get_xtream_player_api_action_url(input, ACTION_GET_SHORT_EPG) { let mut info_url = format!("{action_url}&{TAG_STREAM_ID}={}", pli.provider_id); if !(limit.is_empty() || limit.eq("0")) { info_url = format!("{info_url}&limit={limit}"); @@ -479,7 +481,7 @@ async fn xtream_get_short_epg(app_state: &AppState, user: &ProxyUserCredentials, return HttpResponse::Found().insert_header(("Location", info_url)).finish(); } - return match request_utils::download_text_content(Arc::clone(&app_state.http_client), input, info_url.as_str(), None).await { + return match request::download_text_content(Arc::clone(&app_state.http_client), input, info_url.as_str(), None).await { Ok(content) => HttpResponse::Ok().content_type(mime::APPLICATION_JSON).body(content), Err(err) => { error!("Failed to download epg {}", sanitize_sensitive_info(err.to_string().as_str())); @@ -520,8 +522,8 @@ async fn xtream_get_catchup_response(app_state: &AppState, target: &ConfigTarget let virtual_id: u32 = try_result_bad_request!(FromStr::from_str(stream_id)); let (pli, _) = try_result_bad_request!(xtream_repository::xtream_get_item_for_stream_id(virtual_id, &app_state.config, target, Some(XtreamCluster::Live)).await); let input = try_option_bad_request!(app_state.config.get_input_by_name(pli.input_name.as_str())); - let info_url = try_option_bad_request!(xtream_utils::get_xtream_player_api_action_url(input, ACTION_GET_CATCHUP_TABLE).map(|action_url| format!("{action_url}&{TAG_STREAM_ID}={}&start={start}&end={end}", pli.provider_id))); - let content = try_result_bad_request!(xtream_utils::get_xtream_stream_info_content(Arc::clone(&app_state.http_client), info_url.as_str(), input).await); + let info_url = try_option_bad_request!(xtream::get_xtream_player_api_action_url(input, ACTION_GET_CATCHUP_TABLE).map(|action_url| format!("{action_url}&{TAG_STREAM_ID}={}&start={start}&end={end}", pli.provider_id))); + let content = try_result_bad_request!(xtream::get_xtream_stream_info_content(Arc::clone(&app_state.http_client), info_url.as_str(), input).await); let mut doc: Map = try_result_bad_request!(serde_json::from_str(&content)); let epg_listings = try_option_bad_request!(doc.get_mut(TAG_EPG_LISTINGS).and_then(Value::as_array_mut)); let target_path = try_option_bad_request!(get_target_storage_path(&app_state.config, target.name.as_str())); @@ -583,15 +585,15 @@ async fn xtream_player_api( api_req: UserApiRequest, app_state: &web::Data, ) -> HttpResponse { - let user_target = get_user_target(&api_req, app_state); + let user_target = get_user_target(&api_req, app_state).await; if let Some((user, target)) = user_target { if !target.has_output(&TargetType::Xtream) { - return HttpResponse::Ok().json(get_user_info(&user, &app_state.config)); + return HttpResponse::Ok().json(get_user_info(&user, &app_state.config).await); } let action = api_req.action.trim(); if action.is_empty() { - return HttpResponse::Ok().json(get_user_info(&user, &app_state.config)); + return HttpResponse::Ok().json(get_user_info(&user, &app_state.config).await); } // Process specific playlist actions diff --git a/src/api/main_api.rs b/src/api/main_api.rs index e893ec346..454418114 100644 --- a/src/api/main_api.rs +++ b/src/api/main_api.rs @@ -9,22 +9,22 @@ use std::io::ErrorKind; use std::path::{PathBuf}; use std::sync::Arc; use std::sync::atomic::{AtomicUsize, Ordering}; -use crate::api::hls_api::hls_api_register; -use crate::api::m3u_api::m3u_api_register; +use crate::api::endpoints::hls_api::hls_api_register; +use crate::api::endpoints::m3u_api::m3u_api_register; use crate::api::model::app_state::AppState; use crate::api::model::download::DownloadQueue; -use crate::api::model::shared_stream_manager::SharedStreamManager; +use crate::api::model::streams::shared_stream_manager::SharedStreamManager; use crate::api::scheduler::start_scheduler; -use crate::api::v1_api::v1_api_register; -use crate::api::web_index::index_register; -use crate::api::xmltv_api::xmltv_api_register; -use crate::api::xtream_api::xtream_api_register; +use crate::api::endpoints::v1_api::v1_api_register; +use crate::api::endpoints::web_index::index_register; +use crate::api::endpoints::xmltv_api::xmltv_api_register; +use crate::api::endpoints::xtream_api::xtream_api_register; use crate::model::config::{validate_targets, Config, ProcessTargets, ScheduleConfig}; use crate::model::healthcheck::Healthcheck; -use crate::processing::playlist_processor; -use crate::utils::lru_cache::{LRUResourceCache}; +use crate::processing::processor::playlist; +use crate::tools::lru_cache::{LRUResourceCache}; use crate::utils::size_utils::human_readable_byte_size; -use crate::utils::sys; +use crate::utils::sys_utils; use crate::VERSION; fn get_web_dir_path(web_ui_enabled: bool, web_root: &str) -> Result { @@ -43,7 +43,7 @@ async fn healthcheck(app_state: web::Data,) -> HttpResponse { status: "ok".to_string(), version: VERSION.to_string(), time: ts, - mem: sys::get_memory_usage().map_or(String::from("?"), human_readable_byte_size), + mem: sys_utils::get_memory_usage().map_or(String::from("?"), human_readable_byte_size), active_clients: app_state.active_clients.as_ref().load(Ordering::Relaxed) }) } @@ -81,7 +81,7 @@ fn exec_update_on_boot(client: Arc, cfg: &Arc, targets: let cfg_clone = Arc::clone(cfg); let targets_clone = Arc::clone(targets); actix_rt::spawn( - async move { playlist_processor::exec_processing(client, cfg_clone, targets_clone).await } + async move { playlist::exec_processing(client, cfg_clone, targets_clone).await } ); } } diff --git a/src/api/mod.rs b/src/api/mod.rs index ca4cfc3fe..2169064db 100644 --- a/src/api/mod.rs +++ b/src/api/mod.rs @@ -1,12 +1,5 @@ -pub mod api_utils; -pub mod main_api; -mod download_api; -mod v1_api; -mod xtream_api; -mod m3u_api; -mod xmltv_api; +pub mod model; +mod api_utils; mod scheduler; -mod web_index; - -pub(crate) mod model; -mod hls_api; \ No newline at end of file +mod endpoints; +pub mod main_api; diff --git a/src/api/model/app_state.rs b/src/api/model/app_state.rs index b80570e25..b05941e76 100644 --- a/src/api/model/app_state.rs +++ b/src/api/model/app_state.rs @@ -2,9 +2,9 @@ use std::sync::Arc; use std::sync::atomic::AtomicUsize; use async_std::sync::{Mutex}; use crate::api::model::download::DownloadQueue; -use crate::api::model::shared_stream_manager::SharedStreamManager; +use crate::api::model::streams::shared_stream_manager::SharedStreamManager; use crate::model::config::{Config}; -use crate::utils::lru_cache::LRUResourceCache; +use crate::tools::lru_cache::LRUResourceCache; pub struct AppState { pub config: Arc, diff --git a/src/api/model/mod.rs b/src/api/model/mod.rs index b136434ac..8bb96a2e8 100644 --- a/src/api/model/mod.rs +++ b/src/api/model/mod.rs @@ -1,14 +1,8 @@ -pub mod request; -pub mod config; -pub mod download; -pub mod xtream; pub mod app_state; -pub mod provider_stream; -pub mod persist_pipe_stream; -pub mod provider_stream_factory; -mod buffered_stream; -pub mod shared_stream_manager; -pub mod model_utils; -mod client_stream; -pub mod stream_error; -pub mod active_client_stream; \ No newline at end of file +pub(in crate::api) mod request; +pub(in crate::api) mod config; +pub(in crate::api) mod download; +pub(in crate::api) mod xtream; +pub(in crate::api) mod model_utils; +pub(in crate::api) mod stream_error; +pub(crate) mod streams; diff --git a/src/api/model/model_utils.rs b/src/api/model/model_utils.rs index b505975b1..550d26420 100644 --- a/src/api/model/model_utils.rs +++ b/src/api/model/model_utils.rs @@ -1,10 +1,10 @@ -use crate::debug_if_enabled; +use crate::utils::debug_if_enabled; use actix_web::http::header::{HeaderName, HeaderValue}; use actix_web::{HttpResponseBuilder}; use reqwest::{Response, StatusCode}; use std::collections::{HashSet}; use std::str::FromStr; -use crate::utils::request_utils::sanitize_sensitive_info; +use crate::utils::network::request::sanitize_sensitive_info; const MEDIA_STREAM_HEADERS: &[&str] = &["accept", "content-type", "content-length", "connection", "accept-ranges", "content-range", "vary", "transfer-encoding", "access-control-allow-origin", "access-control-allow-credentials", "icy-metadata"]; diff --git a/src/api/model/active_client_stream.rs b/src/api/model/streams/active_client_stream.rs similarity index 94% rename from src/api/model/active_client_stream.rs rename to src/api/model/streams/active_client_stream.rs index 7c0fa7592..db167377f 100644 --- a/src/api/model/active_client_stream.rs +++ b/src/api/model/streams/active_client_stream.rs @@ -1,4 +1,4 @@ -use crate::api::model::provider_stream_factory::ResponseStream; +use crate::api::model::streams::provider_stream_factory::ResponseStream; use bytes::Bytes; use std::pin::Pin; use std::sync::atomic::{AtomicUsize, Ordering}; diff --git a/src/api/model/buffered_stream.rs b/src/api/model/streams/buffered_stream.rs similarity index 92% rename from src/api/model/buffered_stream.rs rename to src/api/model/streams/buffered_stream.rs index a44bee5a9..baf746392 100644 --- a/src/api/model/buffered_stream.rs +++ b/src/api/model/streams/buffered_stream.rs @@ -1,4 +1,4 @@ -use crate::api::model::provider_stream_factory::ResponseStream; +use crate::api::model::streams::provider_stream_factory::ResponseStream; use futures::{stream::Stream, task::{Context, Poll}, StreamExt}; use std::{ pin::Pin, @@ -7,7 +7,7 @@ use std::{ use tokio::sync::mpsc::channel; use tokio_stream::wrappers::ReceiverStream; use crate::api::model::stream_error::StreamError; -use crate::utils::atomic_once_flag::AtomicOnceFlag; +use crate::tools::atomic_once_flag::AtomicOnceFlag; pub(in crate::api::model) struct BufferedStream { stream: ReceiverStream>, diff --git a/src/api/model/client_stream.rs b/src/api/model/streams/client_stream.rs similarity index 89% rename from src/api/model/client_stream.rs rename to src/api/model/streams/client_stream.rs index 6b8f29c68..d883096d3 100644 --- a/src/api/model/client_stream.rs +++ b/src/api/model/streams/client_stream.rs @@ -1,4 +1,4 @@ -use crate::api::model::provider_stream_factory::ResponseStream; +use crate::api::model::streams::provider_stream_factory::ResponseStream; use bytes::Bytes; use std::pin::Pin; use std::sync::atomic::{AtomicUsize, Ordering}; @@ -6,9 +6,9 @@ use std::sync::{Arc}; use std::task::{Poll}; use futures::{Stream}; use crate::api::model::stream_error::StreamError; -use crate::trace_if_enabled; -use crate::utils::atomic_once_flag::AtomicOnceFlag; -use crate::utils::request_utils::sanitize_sensitive_info; +use crate::utils::trace_if_enabled; +use crate::tools::atomic_once_flag::AtomicOnceFlag; +use crate::utils::network::request::sanitize_sensitive_info; /// This stream counts the send bytes for reconnecting to the actual position and /// sets the `close_signal` if the client drops the connection. diff --git a/src/api/model/streams/mod.rs b/src/api/model/streams/mod.rs new file mode 100644 index 000000000..42f8c712e --- /dev/null +++ b/src/api/model/streams/mod.rs @@ -0,0 +1,7 @@ +pub(in crate::api) mod provider_stream; +pub(in crate::api) mod persist_pipe_stream; +pub(in crate::api) mod provider_stream_factory; +pub(in crate::api) mod shared_stream_manager; +pub(in crate::api) mod active_client_stream; +mod buffered_stream; +mod client_stream; \ No newline at end of file diff --git a/src/api/model/persist_pipe_stream.rs b/src/api/model/streams/persist_pipe_stream.rs similarity index 100% rename from src/api/model/persist_pipe_stream.rs rename to src/api/model/streams/persist_pipe_stream.rs diff --git a/src/api/model/provider_stream.rs b/src/api/model/streams/provider_stream.rs similarity index 92% rename from src/api/model/provider_stream.rs rename to src/api/model/streams/provider_stream.rs index bd1c39a49..d344e8970 100644 --- a/src/api/model/provider_stream.rs +++ b/src/api/model/streams/provider_stream.rs @@ -1,8 +1,7 @@ use crate::api::api_utils::get_headers_from_request; -use crate::api::model::provider_stream_factory::{create_provider_stream, BufferStreamOptions}; -use crate::debug_if_enabled; +use crate::api::model::streams::provider_stream_factory::{create_provider_stream, BufferStreamOptions}; use crate::model::config::ConfigInput; -use crate::utils::request_utils::{get_request_headers, sanitize_sensitive_info}; +use crate::utils::network::request::{get_request_headers, sanitize_sensitive_info}; use actix_web::{HttpRequest}; use bytes::Bytes; use futures::stream::BoxStream; @@ -13,6 +12,7 @@ use futures::TryStreamExt; use url::Url; use crate::api::model::model_utils::get_response_headers; use crate::api::model::stream_error::StreamError; +use crate::utils::debug_if_enabled; type ProviderStreamResponse = (Option>>, Option<(Vec<(String, String)>, StatusCode)>); diff --git a/src/api/model/provider_stream_factory.rs b/src/api/model/streams/provider_stream_factory.rs similarity index 96% rename from src/api/model/provider_stream_factory.rs rename to src/api/model/streams/provider_stream_factory.rs index f3a3f0d67..e3434ad2a 100644 --- a/src/api/model/provider_stream_factory.rs +++ b/src/api/model/streams/provider_stream_factory.rs @@ -1,13 +1,13 @@ use crate::api::api_utils::get_headers_from_request; -use crate::api::model::buffered_stream::BufferedStream; -use crate::api::model::client_stream::ClientStream; +use crate::api::model::streams::buffered_stream::BufferedStream; +use crate::api::model::streams::client_stream::ClientStream; use crate::api::model::model_utils::get_response_headers; use crate::api::model::stream_error::StreamError; -use crate::debug_if_enabled; +use crate::utils::debug_if_enabled; use crate::model::config::ConfigInput; use crate::model::playlist::PlaylistItemType; -use crate::utils::atomic_once_flag::AtomicOnceFlag; -use crate::utils::request_utils::{classify_content_type, get_request_headers, sanitize_sensitive_info, MimeCategory}; +use crate::tools::atomic_once_flag::AtomicOnceFlag; +use crate::utils::network::request::{classify_content_type, get_request_headers, sanitize_sensitive_info, MimeCategory}; use actix_web::HttpRequest; use bytes::Bytes; use futures::stream::{self, BoxStream}; @@ -355,8 +355,8 @@ pub async fn create_provider_stream(client: Arc, #[cfg(test)] mod tests { - use crate::api::model::provider_stream_factory::PlaylistItemType; - use crate::api::model::provider_stream_factory::{create_provider_stream, BufferStreamOptions}; + use crate::api::model::streams::provider_stream_factory::PlaylistItemType; + use crate::api::model::streams::provider_stream_factory::{create_provider_stream, BufferStreamOptions}; use actix_web::test; use actix_web::test::TestRequest; use actix_web::web; diff --git a/src/api/model/shared_stream_manager.rs b/src/api/model/streams/shared_stream_manager.rs similarity index 97% rename from src/api/model/shared_stream_manager.rs rename to src/api/model/streams/shared_stream_manager.rs index ae9194528..deed1464a 100644 --- a/src/api/model/shared_stream_manager.rs +++ b/src/api/model/streams/shared_stream_manager.rs @@ -1,8 +1,8 @@ use crate::api::model::app_state::AppState; -use crate::api::model::provider_stream_factory::STREAM_QUEUE_SIZE; +use crate::api::model::streams::provider_stream_factory::STREAM_QUEUE_SIZE; use crate::api::model::stream_error::StreamError; -use crate::debug_if_enabled; -use crate::utils::request_utils::sanitize_sensitive_info; +use crate::utils::debug_if_enabled; +use crate::utils::network::request::sanitize_sensitive_info; use async_std::sync::Mutex; use bytes::Bytes; use futures::stream::BoxStream; diff --git a/src/api/scheduler.rs b/src/api/scheduler.rs index d14aed530..a066638f0 100644 --- a/src/api/scheduler.rs +++ b/src/api/scheduler.rs @@ -4,9 +4,9 @@ use std::time::{Duration, Instant, SystemTime}; use chrono::{DateTime, FixedOffset, Local}; use cron::Schedule; use log::error; -use crate::exit; +use crate::utils::sys_utils::exit; use crate::model::config::{Config, ProcessTargets}; -use crate::processing::playlist_processor::exec_processing; +use crate::processing::processor::playlist::exec_processing; fn datetime_to_instant(datetime: DateTime) -> Instant { // Convert DateTime to SystemTime diff --git a/src/auth/authenticator.rs b/src/auth/authenticator.rs index 188823fbe..262dd116d 100644 --- a/src/auth/authenticator.rs +++ b/src/auth/authenticator.rs @@ -45,13 +45,15 @@ pub async fn validator( req: ServiceRequest, credentials: Option, ) -> Result { - let app_state: &web::Data = req.app_data::>().unwrap(); - let secret_key = app_state.config.web_auth.as_ref().unwrap().secret.as_ref(); - if verify_token(credentials, secret_key) { - Ok(req) - } else { - Err((actix_web::error::ErrorUnauthorized("Unauthorized"), req)) + if let Some(app_state) = req.app_data::>() { + if let Some(web_auth_config) = app_state.config.web_auth.as_ref() { + let secret_key = web_auth_config.secret.as_ref(); + if verify_token(credentials, secret_key) { + return Ok(req); + } + } } + Err((actix_web::error::ErrorUnauthorized("Unauthorized"), req)) } // pub fn handle_unauthorized(srvres: ServiceResponse) -> actix_web::Result> { diff --git a/src/filter.pest b/src/foundation/filter.pest similarity index 100% rename from src/filter.pest rename to src/foundation/filter.pest diff --git a/src/filter.rs b/src/foundation/filter.rs similarity index 99% rename from src/filter.rs rename to src/foundation/filter.rs index a27b86836..f68d5822e 100644 --- a/src/filter.rs +++ b/src/foundation/filter.rs @@ -12,8 +12,9 @@ use pest::Parser; use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; use crate::model::config::ItemField; use crate::model::playlist::{PlaylistItem, PlaylistItemType}; -use crate::utils::directed_graph::DirectedGraph; -use crate::{create_m3u_filter_error_result, exit, info_err}; +use crate::tools::directed_graph::DirectedGraph; +use crate::m3u_filter_error::{create_m3u_filter_error_result, info_err}; +use crate::utils::sys_utils::exit; pub fn get_field_value(pli: &PlaylistItem, field: &ItemField) -> Rc { let header = pli.header.borrow(); @@ -554,7 +555,7 @@ mod tests { use regex::Regex; - use crate::filter::{get_filter, MockValueProcessor, ValueProvider}; + use crate::foundation::filter::{get_filter, MockValueProcessor, ValueProvider}; use crate::model::playlist::{PlaylistItem, PlaylistItemHeader}; fn create_mock_pli(name: &str, group: &str) -> PlaylistItem { diff --git a/src/foundation/mod.rs b/src/foundation/mod.rs new file mode 100644 index 000000000..ff67a7631 --- /dev/null +++ b/src/foundation/mod.rs @@ -0,0 +1 @@ +pub(crate) mod filter; diff --git a/src/m3u_filter_error.rs b/src/m3u_filter_error.rs index 7f710df5b..f50dfdde3 100644 --- a/src/m3u_filter_error.rs +++ b/src/m3u_filter_error.rs @@ -22,6 +22,8 @@ macro_rules! get_errors_notify_message { }; } +pub use get_errors_notify_message; + #[macro_export] macro_rules! notify_err { ($text:expr) => { @@ -29,12 +31,64 @@ macro_rules! notify_err { }; } +pub use notify_err; + #[macro_export] macro_rules! info_err { ($text:expr) => { M3uFilterError::new(M3uFilterErrorKind::Info, $text) }; } +pub use info_err; + + +#[macro_export] +macro_rules! create_m3u_filter_error { + ($kind: expr, $($arg:tt)*) => { + M3uFilterError::new($kind, format!($($arg)*)) + } +} +pub use create_m3u_filter_error; + +#[macro_export] +macro_rules! create_m3u_filter_error_result { + ($kind: expr, $($arg:tt)*) => { + Err(M3uFilterError::new($kind, format!($($arg)*))) + } +} +pub use create_m3u_filter_error_result; + +#[macro_export] +macro_rules! handle_m3u_filter_error_result_list { + ($kind:expr, $result: expr) => { + let errors = $result + .filter_map(|result| { + if let Err(err) = result { + Some(err.to_string()) + } else { + None + } + }) + .collect::>(); + if !&errors.is_empty() { + return Err(M3uFilterError::new($kind, errors.join("\n"))); + } + } +} + +pub use handle_m3u_filter_error_result_list; + +#[macro_export] +macro_rules! handle_m3u_filter_error_result { + ($kind:expr, $result: expr) => { + if let Err(err) = $result { + return Err(M3uFilterError::new($kind, err.to_string())); + } + } +} + +pub use handle_m3u_filter_error_result; + #[derive(Debug, PartialEq, Eq)] pub enum M3uFilterErrorKind { diff --git a/src/main.rs b/src/main.rs index 09e8e91e6..8d3210126 100644 --- a/src/main.rs +++ b/src/main.rs @@ -1,9 +1,16 @@ +#![warn(clippy::pedantic)] #![allow(clippy::module_name_repetitions)] +#![allow(clippy::must_use_candidate)] +#![allow(clippy::return_self_not_must_use)] +#![allow(clippy::missing_errors_doc)] extern crate core; extern crate env_logger; extern crate pest; #[macro_use] extern crate pest_derive; +#[macro_use] +mod modules; +include_modules!(); use actix_rt::System; use std::fs::File; @@ -13,23 +20,14 @@ use std::sync::Arc; use crate::auth::password::generate_password; use crate::model::config::{validate_targets, Config, HealthcheckConfig, LogLevelConfig, ProcessTargets}; use crate::model::healthcheck::Healthcheck; -use crate::processing::playlist_processor; -use crate::utils::request_utils::set_sanitize_sensitive_info; -use crate::utils::{config_reader, file_utils}; +use crate::processing::processor::playlist; +use crate::utils::config_reader; +use crate::utils::file::file_utils; +use crate::utils::network::request::set_sanitize_sensitive_info; use clap::Parser; use env_logger::Builder; use log::{error, info, LevelFilter}; -mod api; -mod auth; -mod filter; -mod m3u_filter_error; -mod messaging; -mod model; -mod processing; -mod repository; -mod utils; - const LOG_ERROR_LEVEL_MOD: &[&str] = &[ "actix_web::middleware::logger", "reqwest::async_impl::client", @@ -41,6 +39,7 @@ const LOG_ERROR_LEVEL_MOD: &[&str] = &[ "actix_server::accept", ]; + #[derive(Parser)] #[command(name = "m3u-filter")] #[command(author = "euzu ")] @@ -87,6 +86,7 @@ struct Args { healthcheck: bool, } + const VERSION: &str = env!("CARGO_PKG_VERSION"); // #[cfg(not(target_env = "msvc"))] @@ -193,7 +193,7 @@ fn create_directories(cfg: &Config) { fn start_in_cli_mode(cfg: Arc, targets: Arc) { let client = Arc::new(reqwest::Client::new()); - System::new().block_on(async { playlist_processor::exec_processing(client, cfg, targets).await }); + System::new().block_on(async { playlist::exec_processing(client, cfg, targets).await }); } fn start_in_server_mode(cfg: Arc, targets: Arc) { diff --git a/src/model/api_proxy.rs b/src/model/api_proxy.rs index e9e9d3fcd..6e861b823 100644 --- a/src/model/api_proxy.rs +++ b/src/model/api_proxy.rs @@ -4,8 +4,7 @@ use std::str::FromStr; use enum_iterator::Sequence; use log::debug; -use crate::{create_m3u_filter_error_result, info_err}; -use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; +use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind, create_m3u_filter_error_result, info_err}; use crate::utils::config_reader; #[derive(Debug, Clone, serde::Serialize, serde::Deserialize, Sequence, PartialEq, Eq)] @@ -198,7 +197,7 @@ impl ApiProxyConfig { let mut tokens = HashSet::new(); let mut errors = Vec::new(); if self.server.is_empty() { - errors.push("No serverinfo defined".to_string()); + errors.push("No server info defined".to_string()); } else { let mut name_set = HashSet::new(); for server in &self.server { diff --git a/src/model/config.rs b/src/model/config.rs index 14af00879..95bdc6fa5 100644 --- a/src/model/config.rs +++ b/src/model/config.rs @@ -7,25 +7,28 @@ use std::fs::File; use std::io::BufRead; use std::path::PathBuf; use std::str::FromStr; -use std::sync::{Arc, RwLock}; +use std::sync::{Arc}; +use async_std::sync::RwLock; use crate::auth::user::UserCredential; use log::{debug, error, warn}; use path_clean::PathClean; use url::Url; -use crate::filter::{get_filter, prepare_templates, Filter, MockValueProcessor, PatternTemplate, ValueProvider}; +use crate::foundation::filter::{get_filter, prepare_templates, Filter, MockValueProcessor, PatternTemplate, ValueProvider}; +use crate::m3u_filter_error::info_err; use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; use crate::messaging::MsgKind; use crate::model::api_proxy::{ApiProxyConfig, ApiProxyServerInfo, ProxyUserCredentials}; use crate::model::mapping::Mapping; use crate::model::mapping::Mappings; +use crate::utils::config_reader; use crate::utils::default_utils::{default_as_default, default_as_true, default_as_two_u16}; -use crate::utils::file_lock_manager::FileLockManager; -use crate::utils::file_utils::file_reader; +use crate::utils::file::file_lock_manager::FileLockManager; +use crate::utils::file::file_utils; +use crate::utils::file::file_utils::file_reader; use crate::utils::size_utils::parse_size_base_2; -use crate::utils::{config_reader, file_utils}; -use crate::{exit, info_err}; +use crate::utils::sys_utils::exit; pub const MAPPER_ATTRIBUTE_FIELDS: &[&str] = &[ "name", "title", "group", "id", "chno", "logo", @@ -44,47 +47,8 @@ macro_rules! valid_property { $array.contains(&$key) }}; } - -#[macro_export] -macro_rules! create_m3u_filter_error { - ($kind: expr, $($arg:tt)*) => { - M3uFilterError::new($kind, format!($($arg)*)) - } -} - -#[macro_export] -macro_rules! create_m3u_filter_error_result { - ($kind: expr, $($arg:tt)*) => { - Err(M3uFilterError::new($kind, format!($($arg)*))) - } -} - -#[macro_export] -macro_rules! handle_m3u_filter_error_result_list { - ($kind:expr, $result: expr) => { - let errors = $result - .filter_map(|result| { - if let Err(err) = result { - Some(err.to_string()) - } else { - None - } - }) - .collect::>(); - if !&errors.is_empty() { - return Err(M3uFilterError::new($kind, errors.join("\n"))); - } - } -} - -#[macro_export] -macro_rules! handle_m3u_filter_error_result { - ($kind:expr, $result: expr) => { - if let Err(err) = $result { - return Err(M3uFilterError::new($kind, err.to_string())); - } - } -} +pub use valid_property; +use crate::m3u_filter_error::{create_m3u_filter_error_result, handle_m3u_filter_error_result, handle_m3u_filter_error_result_list}; #[derive(Debug, Clone, serde::Serialize, serde::Deserialize, Sequence, PartialEq, Eq, Hash)] pub enum TargetType { @@ -240,12 +204,13 @@ pub struct ConfigSortChannel { impl ConfigSortChannel { pub fn prepare(&mut self) -> Result<(), M3uFilterError> { - let re = regex::Regex::new(&self.group_pattern); - if re.is_err() { - return create_m3u_filter_error_result!(M3uFilterErrorKind::Info, "cant parse regex: {}", &self.group_pattern); + match regex::Regex::new(&self.group_pattern) { + Ok(pattern) => { + self.re = Some(pattern); + Ok(()) + } + Err(err) => create_m3u_filter_error_result!(M3uFilterErrorKind::Info, "cant parse regex: {} {err}", &self.group_pattern), } - self.re = Some(re.unwrap()); - Ok(()) } } @@ -281,12 +246,13 @@ pub struct ConfigRename { impl ConfigRename { pub fn prepare(&mut self) -> Result<(), M3uFilterError> { - let re = regex::Regex::new(&self.pattern); - if re.is_err() { - return create_m3u_filter_error_result!(M3uFilterErrorKind::Info, "cant parse regex: {}", &self.pattern); + match regex::Regex::new(&self.pattern) { + Ok(pattern) => { + self.re = Some(pattern); + Ok(()) + } + Err(err) => create_m3u_filter_error_result!(M3uFilterErrorKind::Info, "cant parse regex: {} {err}", &self.pattern), } - self.re = Some(re.unwrap()); - Ok(()) } } @@ -452,7 +418,10 @@ impl ConfigTarget { pub fn filter(&self, provider: &ValueProvider) -> bool { let mut processor = MockValueProcessor {}; - self.t_filter.as_ref().unwrap().filter(provider, &mut processor) + if let Some(filter) = self.t_filter.as_ref() { + return filter.filter(provider, &mut processor); + } + true } pub fn get_m3u_filename(&self) -> Option<&String> { @@ -637,11 +606,11 @@ impl ConfigInput { pub fn get_user_info(&self) -> Option { if self.input_type == InputType::Xtream { - if self.username.is_some() || self.password.is_some() { + if let (Some(username), Some(password)) = (self.username.as_ref(), self.password.as_ref()) { return Some(InputUserInfo { base_url: self.url.clone(), - username: self.username.as_ref().unwrap().to_owned(), - password: self.password.as_ref().unwrap().to_owned(), + username: username.to_owned(), + password: password.to_owned(), }); } } else if let Ok(url) = Url::parse(&self.url) { @@ -656,11 +625,13 @@ impl ConfigInput { } } if username.is_some() || password.is_some() { - return Some(InputUserInfo { - base_url, - username: username.as_ref().unwrap().to_owned(), - password: password.as_ref().unwrap().to_owned(), - }); + if let (Some(username), Some(password)) = (username.as_ref(), password.as_ref()) { + return Some(InputUserInfo { + base_url, + username: username.to_owned(), + password: password.to_owned(), + }); + } } } None @@ -767,6 +738,10 @@ pub struct VideoConfig { } impl VideoConfig { + + /// # Panics + /// + /// Will panic if default `RegEx` gets invalid pub fn prepare(&mut self) -> Result<(), M3uFilterError> { self.extensions = vec!["mkv".to_string(), "avi".to_string(), "mp4".to_string()]; match &mut self.download { @@ -780,14 +755,16 @@ impl VideoConfig { if let Some(episode_pattern) = &downl.episode_pattern { if !episode_pattern.is_empty() { - let re = regex::Regex::new(episode_pattern); - if re.is_err() { - return create_m3u_filter_error_result!(M3uFilterErrorKind::Info, "cant parse regex: {}", episode_pattern); + match regex::Regex::new(episode_pattern) { + Ok(pattern) => { + downl.t_re_episode_pattern = Some(pattern); + } + Err(err) => { + return create_m3u_filter_error_result!(M3uFilterErrorKind::Info, "cant parse regex: {episode_pattern} {err}"); + } } - downl.t_re_episode_pattern = Some(re.unwrap()); } } - downl.t_re_filename = Some(regex::Regex::new(r"[^A-Za-z0-9_.-]").unwrap()); downl.t_re_remove_filename_ending = Some(regex::Regex::new(r"[_.\s-]$").unwrap()); } @@ -1095,16 +1072,16 @@ impl Config { None } - pub fn get_target_for_user(&self, username: &str, password: &str) -> Option<(ProxyUserCredentials, &ConfigTarget)> { - self.t_api_proxy.read().unwrap().as_ref().and_then(|api_proxy| self.intern_get_target_for_user(api_proxy.get_target_name(username, password))) + pub async fn get_target_for_user(&self, username: &str, password: &str) -> Option<(ProxyUserCredentials, &ConfigTarget)> { + self.t_api_proxy.read().await.as_ref().and_then(|api_proxy| self.intern_get_target_for_user(api_proxy.get_target_name(username, password))) } - pub fn get_target_for_user_by_token(&self, token: &str) -> Option<(ProxyUserCredentials, &ConfigTarget)> { - self.t_api_proxy.read().unwrap().as_ref().and_then(|api_proxy| self.intern_get_target_for_user(api_proxy.get_target_name_by_token(token))) + pub async fn get_target_for_user_by_token(&self, token: &str) -> Option<(ProxyUserCredentials, &ConfigTarget)> { + self.t_api_proxy.read().await.as_ref().and_then(|api_proxy| self.intern_get_target_for_user(api_proxy.get_target_name_by_token(token))) } - pub fn get_user_credentials(&self, username: &str) -> Option { - self.t_api_proxy.read().unwrap().as_ref().and_then(|api_proxy| api_proxy.get_user_credentials(username)) + pub async fn get_user_credentials(&self, username: &str) -> Option { + self.t_api_proxy.read().await.as_ref().and_then(|api_proxy| api_proxy.get_user_credentials(username)) } pub fn get_input_by_name(&self, input_name: &str) -> Option<&ConfigInput> { @@ -1192,11 +1169,17 @@ impl Config { pub fn prepare(&mut self, resolve_var: bool) -> Result<(), M3uFilterError> { let work_dir = if resolve_var { &config_reader::resolve_env_var(&self.working_dir) } else { &self.working_dir }; self.working_dir = file_utils::get_working_path(work_dir); + if self.backup_dir.is_none() { self.backup_dir = Some(PathBuf::from(&self.working_dir).join("backup").clean().to_string_lossy().to_string()); } else { - let backup_dir = if resolve_var { &config_reader::resolve_env_var(self.backup_dir.as_ref().unwrap()) } else { self.backup_dir.as_ref().unwrap() }; - self.backup_dir = Some(backup_dir.to_string()); + self.backup_dir = self.backup_dir.as_ref().map(|backup_dir| { + if resolve_var { + config_reader::resolve_env_var(backup_dir) + } else { + backup_dir.to_owned() + } + }).map(|dir| dir.to_string()); } if let Some(reverse_proxy) = self.reverse_proxy.as_mut() { reverse_proxy.prepare(&self.working_dir, resolve_var); @@ -1286,8 +1269,12 @@ impl Config { } } - pub fn get_user_server_info(&self, user: &ProxyUserCredentials) -> ApiProxyServerInfo { - let server_info_list = self.t_api_proxy.read().unwrap().as_ref().unwrap().server.clone(); + + /// # Panics + /// + /// Will panic if default server invalid + pub async fn get_user_server_info(&self, user: &ProxyUserCredentials) -> ApiProxyServerInfo { + let server_info_list = self.t_api_proxy.read().await.as_ref().unwrap().server.clone(); let server_info_name = user.server.as_ref().map_or("default", |server_name| server_name.as_str()); server_info_list.iter().find(|c| c.name.eq(server_info_name)).map_or_else(|| server_info_list.first().unwrap().clone(), Clone::clone) } diff --git a/src/model/mapping.rs b/src/model/mapping.rs index e67a41848..098745791 100644 --- a/src/model/mapping.rs +++ b/src/model/mapping.rs @@ -1,21 +1,22 @@ +use enum_iterator::Sequence; +use log::{debug, error, trace}; +use regex::{Regex}; use std::borrow::Cow; use std::cell::RefCell; use std::collections::HashMap; use std::fmt::Display; use std::rc::Rc; use std::str::FromStr; -use std::sync::{Arc}; use std::sync::atomic::AtomicU32; -use enum_iterator::Sequence; -use log::{debug, error, trace}; -use regex::Regex; +use std::sync::Arc; -use crate::filter::{apply_templates_to_pattern, get_filter, prepare_templates, Filter, PatternTemplate, RegexWithCaptures, ValueProcessor}; +use crate::foundation::filter::{apply_templates_to_pattern, get_filter, prepare_templates, Filter, PatternTemplate, RegexWithCaptures, ValueProcessor}; +use crate::m3u_filter_error::{create_m3u_filter_error_result, handle_m3u_filter_error_result, info_err}; use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; +use crate::model::config::valid_property; use crate::model::config::{ItemField, AFFIX_FIELDS, COUNTER_FIELDS, MAPPER_ATTRIBUTE_FIELDS}; use crate::model::playlist::{FieldGetAccessor, FieldSetAccessor, PlaylistItem}; use crate::utils::string_utils::Capitalize; -use crate::{create_m3u_filter_error_result, handle_m3u_filter_error_result, info_err, valid_property}; #[derive(Debug, Clone, serde::Serialize, serde::Deserialize, Default)] pub struct MappingTag { @@ -168,11 +169,10 @@ impl MapperTransform { new_pattern = apply_templates_to_pattern(pattern, template_list); } } - let re = Regex::new(&new_pattern); - if re.is_err() { - return create_m3u_filter_error_result!(M3uFilterErrorKind::Info, "cant parse regex: {}", new_pattern); + match Regex::new(&new_pattern) { + Ok(pattern) => self.t_pattern = Some(pattern), + Err(err) => return create_m3u_filter_error_result!(M3uFilterErrorKind::Info, "cant parse regex: {new_pattern} {err}"), } - self.t_pattern = Some(re.unwrap()); } } Ok(()) @@ -206,6 +206,9 @@ pub struct Mapper { } impl Mapper { + /// # Panics + /// + /// Will panic if default `RegEx` gets invalid pub fn prepare(&mut self, templates: Option<&Vec>, tags: Option<&Vec>) -> Result<(), M3uFilterError> { for key in self.attributes.keys() { if !valid_property!(key.as_str(), MAPPER_ATTRIBUTE_FIELDS) { diff --git a/src/model/playlist.rs b/src/model/playlist.rs index 9eb516d8d..1e88bad32 100644 --- a/src/model/playlist.rs +++ b/src/model/playlist.rs @@ -11,7 +11,7 @@ use crate::utils::json_utils::{get_string_from_serde_value, get_u64_from_serde_v use serde::{Deserialize, Serialize}; use serde_json::{Map, Value}; use crate::utils::hash_utils::{generate_playlist_uuid, get_provider_id}; -use crate::utils::request_utils::extract_extension_from_url; +use crate::utils::network::request::extract_extension_from_url; // https://de.wikipedia.org/wiki/M3U // https://siptv.eu/howto/playlist.html diff --git a/src/model/xtream.rs b/src/model/xtream.rs index f7da7c0f2..4739ee8bb 100644 --- a/src/model/xtream.rs +++ b/src/model/xtream.rs @@ -258,7 +258,7 @@ impl XtreamStream { let mut result = Map::new(); if let Some(bdpath) = self.backdrop_path.as_ref() { if !bdpath.is_empty() { - result.insert(String::from(PROP_BACKDROP_PATH), Value::Array(Vec::from([Value::String(String::from(bdpath.first().unwrap()))]))); + result.insert(String::from(PROP_BACKDROP_PATH), Value::Array(Vec::from([Value::String(String::from(bdpath.first()?))]))); } } add_rc_str_property_if_exists!(result, self.tmdb, "tmdb"); @@ -461,7 +461,7 @@ impl XtreamSeriesInfoEpisode { let bdpath = info.and_then(|i| i.backdrop_path.as_ref()); let bdpath_is_set = bdpath.as_ref().is_some_and(|bdpath| !bdpath.is_empty()); if bdpath_is_set { - result.insert(String::from("backdrop_path"), Value::Array(Vec::from([Value::String(String::from(bdpath.as_ref().unwrap().first()?))]))); + result.insert(String::from("backdrop_path"), Value::Array(Vec::from([Value::String(String::from(bdpath?.first()?))]))); } add_str_property_if_exists!(result, info.map_or("", |i| i.name.as_str()), "series_name"); add_str_property_if_exists!(result, info.map_or("", |i| i.release_date.as_str()), "series_release_date"); diff --git a/src/modules.rs b/src/modules.rs new file mode 100644 index 000000000..25d453511 --- /dev/null +++ b/src/modules.rs @@ -0,0 +1,16 @@ +#[macro_export] +macro_rules! include_modules { + () => { + pub mod api; + pub mod auth; + pub mod m3u_filter_error; + pub mod messaging; + pub mod model; + pub mod processing; + pub mod repository; + pub mod utils; + pub mod tools; + pub mod foundation; + } +} + diff --git a/src/processing/mod.rs b/src/processing/mod.rs index 0b67a5980..e1ff80590 100644 --- a/src/processing/mod.rs +++ b/src/processing/mod.rs @@ -1,10 +1,3 @@ -pub mod m3u_parser; -pub mod xtream_parser; -pub mod playlist_processor; -pub mod xmltv_parser; mod playlist_watch; -mod xtream_processor; -mod affix_processor; -mod xtream_processor_vod; -mod xtream_processor_series; -pub mod hls_parser; +pub(crate) mod parser; +pub(crate) mod processor; diff --git a/src/processing/hls_parser.rs b/src/processing/parser/hls.rs similarity index 100% rename from src/processing/hls_parser.rs rename to src/processing/parser/hls.rs diff --git a/src/processing/m3u_parser.rs b/src/processing/parser/m3u.rs similarity index 100% rename from src/processing/m3u_parser.rs rename to src/processing/parser/m3u.rs diff --git a/src/processing/parser/mod.rs b/src/processing/parser/mod.rs new file mode 100644 index 000000000..5cfbaf9b8 --- /dev/null +++ b/src/processing/parser/mod.rs @@ -0,0 +1,4 @@ +pub mod m3u; +pub mod xtream; +pub mod xmltv; +pub mod hls; \ No newline at end of file diff --git a/src/processing/xmltv_parser.rs b/src/processing/parser/xmltv.rs similarity index 98% rename from src/processing/xmltv_parser.rs rename to src/processing/parser/xmltv.rs index 487f87a72..9787cff43 100644 --- a/src/processing/xmltv_parser.rs +++ b/src/processing/parser/xmltv.rs @@ -5,7 +5,7 @@ use quick_xml::events::Event; use quick_xml::Reader; use crate::model::xmltv::{Epg, EPG_ATTRIB_CHANNEL, EPG_ATTRIB_ID, EPG_TAG_TV, EPG_TAG_CHANNEL, EPG_TAG_PROGRAMME, TVGuide, XmlTag}; -use crate::utils::compressed_file_reader::CompressedFileReader; +use crate::utils::compression::compressed_file_reader::CompressedFileReader; impl TVGuide { pub fn filter(&self, channel_ids: &HashSet>) -> Option { diff --git a/src/processing/xtream_parser.rs b/src/processing/parser/xtream.rs similarity index 98% rename from src/processing/xtream_parser.rs rename to src/processing/parser/xtream.rs index cccb52d9a..6087d36e0 100644 --- a/src/processing/xtream_parser.rs +++ b/src/processing/parser/xtream.rs @@ -4,13 +4,12 @@ use std::rc::Rc; use serde_json::Value; -use crate::create_m3u_filter_error_result; -use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; +use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind, create_m3u_filter_error_result}; use crate::model::config::ConfigInput; use crate::model::playlist::{PlaylistGroup, PlaylistItem, PlaylistItemHeader, PlaylistItemType, XtreamCluster}; use crate::model::xtream::{XtreamCategory, XtreamSeriesInfo, XtreamSeriesInfoEpisode, XtreamStream}; use crate::utils::hash_utils::generate_playlist_uuid; -use crate::utils::xtream_utils::{get_xtream_stream_url_base, ACTION_GET_SERIES_INFO}; +use crate::utils::network::xtream::{get_xtream_stream_url_base, ACTION_GET_SERIES_INFO}; fn map_to_xtream_category(categories: &Value) -> Result, M3uFilterError> { match serde_json::from_value::>(categories.to_owned()) { diff --git a/src/processing/playlist_watch.rs b/src/processing/playlist_watch.rs index ced59b548..ef327c7fe 100644 --- a/src/processing/playlist_watch.rs +++ b/src/processing/playlist_watch.rs @@ -4,8 +4,8 @@ use log::{error, info}; use crate::messaging::{MsgKind, send_message}; use crate::model::config::Config; use crate::model::playlist::PlaylistGroup; -use crate::utils::file_utils; -use crate::utils::file_utils::sanitize_filename; +use crate::utils::file::file_utils; +use crate::utils::file::file_utils::sanitize_filename; pub fn process_group_watch(cfg: &Config, target_name: &str, pl: &PlaylistGroup) { let mut new_tree = BTreeSet::new(); diff --git a/src/processing/affix_processor.rs b/src/processing/processor/affix.rs similarity index 97% rename from src/processing/affix_processor.rs rename to src/processing/processor/affix.rs index 707459f8c..ef633d04e 100644 --- a/src/processing/affix_processor.rs +++ b/src/processing/processor/affix.rs @@ -1,6 +1,6 @@ -use crate::model::config::{ConfigInput, InputAffix, AFFIX_FIELDS}; +use crate::model::config::{ConfigInput, InputAffix, AFFIX_FIELDS, valid_property}; use crate::model::playlist::{FetchedPlaylist, FieldGetAccessor, FieldSetAccessor, PlaylistItem}; -use crate::{debug_if_enabled, valid_property}; +use crate::utils::{debug_if_enabled}; type AffixProcessor<'a> = Box; diff --git a/src/processing/processor/mod.rs b/src/processing/processor/mod.rs new file mode 100644 index 000000000..768fe116c --- /dev/null +++ b/src/processing/processor/mod.rs @@ -0,0 +1,44 @@ +pub mod playlist; +mod xtream; +mod affix; +mod xtream_vod; +mod xtream_series; + +#[macro_export] +macro_rules! handle_error { + ($stmt:expr, $map_err:expr) => { + if let Err(err) = $stmt { + $map_err(err); + } + }; +} +use handle_error; + +#[macro_export] +macro_rules! handle_error_and_return { + ($stmt:expr, $map_err:expr) => { + if let Err(err) = $stmt { + $map_err(err); + return Default::default(); + } + }; +} +use handle_error_and_return; + + +#[macro_export] +macro_rules! create_resolve_options_function_for_xtream_target { + ($cluster:ident) => { + paste::paste! { + fn [](target: &ConfigTarget, fpl: &FetchedPlaylist) -> (bool, u16) { + let (resolve, resolve_delay) = + target.options.as_ref().map_or((false, 0), |opt| { + (opt.[] && fpl.input.input_type == InputType::Xtream, + opt.[]) + }); + (resolve, resolve_delay) + } + } + }; +} +use create_resolve_options_function_for_xtream_target; diff --git a/src/processing/playlist_processor.rs b/src/processing/processor/playlist.rs similarity index 96% rename from src/processing/playlist_processor.rs rename to src/processing/processor/playlist.rs index de8dfe911..c6de08dae 100644 --- a/src/processing/playlist_processor.rs +++ b/src/processing/processor/playlist.rs @@ -1,8 +1,10 @@ extern crate unidecode; -use crate::utils::epg_utils; -use crate::utils::m3u_utils; -use crate::utils::xtream_utils; +use crate::Config; +use crate::model::config::ConfigRename; +use crate::utils::network::epg; +use crate::utils::network::m3u; +use crate::utils::network::xtream; use async_std::sync::Mutex; use core::cmp::Ordering; use std::cell::RefCell; @@ -17,22 +19,22 @@ use log::{debug, error, info, log_enabled, trace, warn, Level}; use std::time::Instant; use unidecode::unidecode; -use crate::filter::{get_field_value, set_field_value, MockValueProcessor, ValueProvider}; -use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; +use crate::foundation::filter::{get_field_value, set_field_value, MockValueProcessor, ValueProvider}; +use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind, get_errors_notify_message, notify_err}; use crate::messaging::{send_message, MsgKind}; use crate::model::config::{ConfigSortChannel, ConfigSortGroup, ConfigTarget, InputType, ItemField, ProcessTargets, ProcessingOrder, SortOrder::{Asc, Desc}}; use crate::model::mapping::{CounterModifier, Mapping, MappingValueProcessor}; use crate::model::playlist::{FetchedPlaylist, FieldGetAccessor, FieldSetAccessor, PlaylistEntry, PlaylistGroup, PlaylistItem, UUIDType, XtreamCluster}; use crate::model::stats::{InputStats, PlaylistStats, SourceStats, TargetStats}; -use crate::processing::affix_processor::apply_affixes; +use crate::processing::processor::affix::apply_affixes; use crate::processing::playlist_watch::process_group_watch; -use crate::processing::xmltv_parser::flatten_tvguide; -use crate::processing::xtream_processor_series::playlist_resolve_series; -use crate::processing::xtream_processor_vod::playlist_resolve_vod; +use crate::processing::parser::xmltv::flatten_tvguide; +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::default_utils::default_as_default; -use crate::{debug_if_enabled, get_errors_notify_message, model::config, notify_err, Config}; +use crate::utils::{debug_if_enabled}; fn is_valid(pli: &PlaylistItem, target: &ConfigTarget) -> bool { let provider = ValueProvider { pli: RefCell::new(pli) }; @@ -140,7 +142,7 @@ fn channel_no_playlist(new_playlist: &[PlaylistGroup]) { } } -fn exec_rename(pli: &PlaylistItem, rename: Option<&Vec>) { +fn exec_rename(pli: &PlaylistItem, rename: Option<&Vec>) { if let Some(renames) = rename { if !renames.is_empty() { let result = pli; @@ -313,11 +315,11 @@ async fn process_source(client: Arc, cfg: Arc, source_i if is_input_enabled(enabled_inputs, input.enabled, input.id, &user_targets) { let start_time = Instant::now(); let (mut playlistgroups, mut error_list) = match input.input_type { - InputType::M3u => m3u_utils::get_m3u_playlist(Arc::clone(&client), &cfg, input, &cfg.working_dir).await, - InputType::Xtream => xtream_utils::get_xtream_playlist(Arc::clone(&client), input, &cfg.working_dir).await, + InputType::M3u => m3u::get_m3u_playlist(Arc::clone(&client), &cfg, input, &cfg.working_dir).await, + InputType::Xtream => xtream::get_xtream_playlist(Arc::clone(&client), input, &cfg.working_dir).await, }; let (tvguide, mut tvguide_errors) = if error_list.is_empty() { - epg_utils::get_xmltv(Arc::clone(&client), &cfg, input, &cfg.working_dir).await + epg::get_xmltv(Arc::clone(&client), &cfg, input, &cfg.working_dir).await } else { (None, vec![]) }; diff --git a/src/processing/xtream_processor.rs b/src/processing/processor/xtream.rs similarity index 79% rename from src/processing/xtream_processor.rs rename to src/processing/processor/xtream.rs index e2809ee98..44fb7b388 100644 --- a/src/processing/xtream_processor.rs +++ b/src/processing/processor/xtream.rs @@ -2,8 +2,7 @@ use crate::m3u_filter_error::{str_to_io_error, to_io_error, M3uFilterError, M3uF use crate::model::config::{Config, ConfigInput}; use crate::model::playlist::{FetchedPlaylist, PlaylistEntry, PlaylistItem, PlaylistItemType, XtreamCluster}; use crate::repository::storage::get_input_storage_path; -use crate::utils::{xtream_utils}; -use crate::{info_err, notify_err}; +use crate::m3u_filter_error::{info_err, notify_err}; use serde::{Deserialize, Serialize}; use std::collections::HashMap; use std::fs::File; @@ -16,49 +15,16 @@ const FILE_VOD_INFO: &str = "xtream_vod_info"; const FILE_SUFFIX_WAL: &str = "wal"; const FILE_SERIES_EPISODE_RECORD: &str = "series_episode_record"; -#[macro_export] -macro_rules! handle_error { - ($stmt:expr, $map_err:expr) => { - if let Err(err) = $stmt { - $map_err(err); - } - }; -} - -#[macro_export] -macro_rules! handle_error_and_return { - ($stmt:expr, $map_err:expr) => { - if let Err(err) = $stmt { - $map_err(err); - return Default::default(); - } - }; -} use crate::repository::bplustree::BPlusTree; use crate::repository::xtream_repository::xtream_get_record_file_path; -use crate::utils::file_utils::append_or_crate_file; - -#[macro_export] -macro_rules! create_resolve_options_function_for_xtream_target { - ($cluster:ident) => { - paste::paste! { - fn [](target: &ConfigTarget, fpl: &FetchedPlaylist) -> (bool, u16) { - let (resolve, resolve_delay) = - target.options.as_ref().map_or((false, 0), |opt| { - (opt.[] && fpl.input.input_type == InputType::Xtream, - opt.[]) - }); - (resolve, resolve_delay) - } - } - }; -} +use crate::utils::file::file_utils::append_or_crate_file; +use crate::utils::network::xtream; pub(in crate::processing) async fn playlist_resolve_download_playlist_item(client: Arc, pli: &PlaylistItem, input: &ConfigInput, errors: &mut Vec, resolve_delay: u16, cluster: XtreamCluster) -> Option { let mut result = None; let provider_id = pli.get_provider_id()?; - if let Some(info_url) = xtream_utils::get_xtream_player_api_info_url(input, cluster, provider_id) { - result = match xtream_utils::get_xtream_stream_info_content(client, &info_url, input).await { + if let Some(info_url) = xtream::get_xtream_player_api_info_url(input, cluster, provider_id) { + result = match xtream::get_xtream_stream_info_content(client, &info_url, input).await { Ok(content) => Some(content), Err(err) => { errors.push(info_err!(format!("{err}"))); diff --git a/src/processing/xtream_processor_series.rs b/src/processing/processor/xtream_series.rs similarity index 95% rename from src/processing/xtream_processor_series.rs rename to src/processing/processor/xtream_series.rs index 25ee35312..c7dff3222 100644 --- a/src/processing/xtream_processor_series.rs +++ b/src/processing/processor/xtream_series.rs @@ -1,13 +1,14 @@ use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; use crate::model::config::{Config, ConfigTarget, InputType}; use crate::model::playlist::{FetchedPlaylist, PlaylistGroup, PlaylistItem, PlaylistItemType, XtreamCluster}; -use crate::processing::playlist_processor::ProcessingPipe; -use crate::processing::xtream_parser::parse_xtream_series_info; -use crate::processing::xtream_processor::{create_resolve_episode_wal_files, create_resolve_info_wal_files, playlist_resolve_download_playlist_item, read_processed_info_ids, should_update_info, write_info_content_to_wal_file}; +use crate::processing::processor::playlist::ProcessingPipe; +use crate::processing::parser::xtream::parse_xtream_series_info; +use crate::processing::processor::xtream::{create_resolve_episode_wal_files, create_resolve_info_wal_files, playlist_resolve_download_playlist_item, read_processed_info_ids, should_update_info, write_info_content_to_wal_file}; use crate::repository::storage::get_input_storage_path; use crate::repository::xtream_repository::{xtream_get_info_file_paths, xtream_update_input_info_file, xtream_update_input_series_episodes_record_from_wal_file, xtream_update_input_series_record_from_wal_file}; use crate::repository::IndexedDocumentReader; -use crate::{create_resolve_options_function_for_xtream_target, handle_error, handle_error_and_return, info_err, notify_err}; +use crate::m3u_filter_error::{notify_err, info_err}; +use crate::processing::processor::{handle_error, handle_error_and_return, create_resolve_options_function_for_xtream_target}; use std::collections::HashMap; use std::fs::File; use std::io::{BufWriter, Write}; @@ -15,7 +16,7 @@ use std::sync::Arc; use std::time::Instant; use log::{info, log_enabled, Level}; use crate::model::xtream::{XtreamSeriesEpisode, XtreamSeriesInfoEpisode}; -use crate::utils::file_utils::file_writer; +use crate::utils::file::file_utils::file_writer; const TAG_SERIES_INFO_LAST_MODIFIED: &str = "last_modified"; diff --git a/src/processing/xtream_processor_vod.rs b/src/processing/processor/xtream_vod.rs similarity index 93% rename from src/processing/xtream_processor_vod.rs rename to src/processing/processor/xtream_vod.rs index c2a548f18..8acc10a97 100644 --- a/src/processing/xtream_processor_vod.rs +++ b/src/processing/processor/xtream_vod.rs @@ -1,9 +1,10 @@ use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; use crate::model::config::{Config, ConfigTarget, InputType}; use crate::model::playlist::{FetchedPlaylist, PlaylistItem, PlaylistItemType, XtreamCluster}; -use crate::processing::xtream_processor::{create_resolve_info_wal_files, playlist_resolve_download_playlist_item, read_processed_info_ids, should_update_info, write_info_content_to_wal_file}; +use crate::processing::processor::xtream::{create_resolve_info_wal_files, playlist_resolve_download_playlist_item, read_processed_info_ids, should_update_info, write_info_content_to_wal_file}; use crate::repository::xtream_repository::{xtream_update_input_info_file, xtream_update_input_vod_record_from_wal_file, InputVodInfoRecord}; -use crate::{create_resolve_options_function_for_xtream_target, handle_error, handle_error_and_return, notify_err}; +use crate::m3u_filter_error::{notify_err}; +use crate::processing::processor::{handle_error, handle_error_and_return, create_resolve_options_function_for_xtream_target}; use crate::utils::json_utils::{get_u32_from_serde_value, get_u64_from_serde_value}; use serde_json::{Map, Value}; use std::collections::HashMap; @@ -12,7 +13,7 @@ use std::io::{BufWriter, Write}; use std::sync::Arc; use std::time::Instant; use log::{info, log_enabled, Level}; -use crate::utils::file_utils::file_writer; +use crate::utils::file::file_utils::file_writer; const TAG_VOD_INFO_INFO: &str = "info"; const TAG_VOD_INFO_MOVIE_DATA: &str = "movie_data"; diff --git a/src/repository/bplustree.rs b/src/repository/bplustree.rs index f317dae62..528df9dc5 100644 --- a/src/repository/bplustree.rs +++ b/src/repository/bplustree.rs @@ -1,16 +1,16 @@ -use std::fs::{File}; +use std::fs::File; use std::io::{self, BufReader, Read, Seek, SeekFrom, Write}; use std::marker::PhantomData; use std::mem::size_of; use std::path::Path; +use crate::m3u_filter_error::{str_to_io_error, to_io_error}; +use crate::utils::file::file_utils::{file_reader, file_writer, open_read_write_file, rename_or_copy}; use log::error; use ruzstd::decoding::StreamingDecoder; use ruzstd::encoding::{compress_to_vec, CompressionLevel}; use serde::{Deserialize, Serialize}; use tempfile::NamedTempFile; -use crate::m3u_filter_error::{str_to_io_error, to_io_error}; -use crate::utils::file_utils::{file_reader, file_writer, open_read_write_file, rename_or_copy}; const BLOCK_SIZE: usize = 4096; const BINCODE_OVERHEAD: usize = 8; @@ -115,9 +115,9 @@ where order >> 1 } - fn find_leaf_entry(node: &Self) -> &K { + fn find_leaf_entry(node: &Self) -> Option<&K> { if node.is_leaf { - node.keys.first().unwrap() + node.keys.first() } else { let child = node.children.first().unwrap(); Self::find_leaf_entry(child) @@ -176,13 +176,14 @@ where let child = self.children.get_mut(pos)?; let node = child.insert(key.clone(), v, inner_order, leaf_order); if node.is_some() { - let leaf_key = Self::find_leaf_entry(node.as_ref().unwrap()); - let idx = self.get_entry_index_upper_bound(leaf_key); - if self.keys.binary_search(&key).is_err() { - self.keys.insert(idx, leaf_key.clone()); - self.children.insert(idx + 1, node.unwrap()); - if self.is_overflow(inner_order) { - return Some(self.split(inner_order)); + if let Some(leaf_key) = Self::find_leaf_entry(node.as_ref().unwrap()) { + let idx = self.get_entry_index_upper_bound(leaf_key); + if self.keys.binary_search(&key).is_err() { + self.keys.insert(idx, leaf_key.clone()); + self.children.insert(idx + 1, node.unwrap()); + if self.is_overflow(inner_order) { + return Some(self.split(inner_order)); + } } } } @@ -329,7 +330,7 @@ where left_over_bytes -= bytes_available_on_block; while left_over_bytes > 0 { file.read_exact(buffer)?; - let bytes_to_read = if left_over_bytes >= BLOCK_SIZE { + let bytes_to_read = if left_over_bytes >= BLOCK_SIZE { BLOCK_SIZE } else { left_over_bytes @@ -366,7 +367,7 @@ where let nodes: Result, io::Error> = pointers .iter() .map(|pointer| { - Self::deserialize_from_block(file, buffer, *pointer, nested) + Self::deserialize_from_block(file, buffer, *pointer, nested) .map(|(node, _)| node) }) .collect(); @@ -385,7 +386,7 @@ fn decode_content(content_bytes: &Vec) -> Option> { if let Ok(mut decoder) = StreamingDecoder::new(&**content_bytes) { let mut result = Vec::with_capacity(content_bytes.len()); if decoder.read_to_end(&mut result).is_ok() { - return Some(result) + return Some(result); } } @@ -408,7 +409,7 @@ pub struct BPlusTree { } const fn calc_order() -> (usize, usize) { - let overhead_size = BINCODE_OVERHEAD + LEN_SIZE + FLAG_SIZE; + let overhead_size = BINCODE_OVERHEAD + LEN_SIZE + FLAG_SIZE; let key_size = size_of::() + overhead_size; let value_size = key_size + size_of::() + overhead_size; let inner_order = BLOCK_SIZE / key_size; @@ -416,6 +417,16 @@ const fn calc_order() -> (usize, usize) { (inner_order, leaf_order) } +impl Default for BPlusTree +where + K: Ord + Serialize + for<'de> Deserialize<'de> + Clone, + V: Serialize + for<'de> Deserialize<'de> + Clone, +{ + fn default() -> Self { + Self::new() + } +} + impl BPlusTree where K: Ord + Serialize + for<'de> Deserialize<'de> + Clone, @@ -450,18 +461,22 @@ where } if let Some(node) = self.root.insert(key, value, self.inner_order, self.leaf_order) { - let child_key = if node.is_leaf { - node.keys.first().as_ref().unwrap() + let child_key_opt = if node.is_leaf { + node.keys.first() } else { BPlusTreeNode::::find_leaf_entry(&node) }; - let mut new_root = BPlusTreeNode::::new(false); - new_root.keys.push(child_key.clone()); - new_root.children.push(std::mem::replace(&mut self.root, BPlusTreeNode::new(true))); - new_root.children.push(node); + if let Some(child_key) = child_key_opt { + let mut new_root = BPlusTreeNode::::new(false); + new_root.keys.push(child_key.clone()); + new_root.children.push(std::mem::replace(&mut self.root, BPlusTreeNode::new(true))); + new_root.children.push(node); - self.root = new_root; + self.root = new_root; + } else { + error!("Failed to insert child key"); + } } } @@ -664,7 +679,11 @@ where }; } let child_idx = get_entry_index_upper_bound::(&node.keys, key); - offset = *pointers.unwrap().get(child_idx).unwrap(); + if let Some(pters) = pointers { + if let Some(child_idx) = pters.get(child_idx) { + offset = *child_idx; + } + } } Err(err) => { error!("Failed to read id tree from file {err}"); @@ -744,11 +763,24 @@ where K: Ord + Serialize + for<'de> Deserialize<'de> + Clone, V: Serialize + for<'de> Deserialize<'de> + Clone, { - pub fn iter(&self) -> BPlusTreeIterator { + pub fn iter(&self) -> BPlusTreeIterator<'_, K, V> { BPlusTreeIterator::new(self) } } +impl<'a, K, V> IntoIterator for &'a BPlusTree +where + K: Ord + Serialize + for<'de> Deserialize<'de> + Clone, + V: Serialize + for<'de> Deserialize<'de> + Clone, +{ + type Item = (&'a K, &'a V); + type IntoIter = BPlusTreeIterator<'a, K, V>; + + fn into_iter(self) -> Self::IntoIter { + self.iter() + } +} + #[cfg(test)] mod tests { use std::collections::HashSet; @@ -785,7 +817,7 @@ mod tests { #[test] fn insert_test() -> io::Result<()> { let test_size = 500; - let content = generate_random_string(1024); + let content = generate_random_string(1024); let mut tree = BPlusTree::::new(); for i in 0u32..=test_size { tree.insert(i, Record { diff --git a/src/repository/epg_repository.rs b/src/repository/epg_repository.rs index 69b4695c0..b591c4dd5 100644 --- a/src/repository/epg_repository.rs +++ b/src/repository/epg_repository.rs @@ -2,8 +2,8 @@ use std::fs::File; use std::io::{Cursor, Write}; use std::path::{Path}; use quick_xml::{Writer}; -use crate::{debug_if_enabled, notify_err}; -use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; +use crate::utils::{debug_if_enabled}; +use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind, notify_err}; use crate::model::config::{Config, ConfigTarget, TargetOutput}; use crate::model::config::TargetType; use crate::model::xmltv::{Epg}; diff --git a/src/repository/indexed_document.rs b/src/repository/indexed_document.rs index b7906023e..e41bd1bfb 100644 --- a/src/repository/indexed_document.rs +++ b/src/repository/indexed_document.rs @@ -5,12 +5,12 @@ use std::marker::PhantomData; use std::path::{Path, PathBuf}; use crate::repository::bplustree::{BPlusTree, BPlusTreeQuery}; -use crate::utils::file_utils; +use crate::utils::file::file_utils; use log::error; use serde::{Deserialize, Serialize}; use tempfile::NamedTempFile; use crate::m3u_filter_error::{str_to_io_error, to_io_error}; -use crate::utils::file_utils::{create_new_file_for_read_write, file_reader, file_writer, open_read_write_file, open_readonly_file, rename_or_copy}; +use crate::utils::file::file_utils::{create_new_file_for_read_write, file_reader, file_writer, open_read_write_file, open_readonly_file, rename_or_copy}; const BLOCK_SIZE: usize = 4096; const LEN_SIZE: usize = 4; @@ -257,7 +257,7 @@ where t_type: PhantomData, }) } else { - Err(Error::new(ErrorKind::NotFound, format!("File not found {}", main_path.to_str().unwrap()))) + Err(Error::new(ErrorKind::NotFound, format!("File not found {main_path:?}"))) } } pub fn get(&mut self, doc_id: &K) -> Result { diff --git a/src/repository/kodi_repository.rs b/src/repository/kodi_repository.rs index 6b5997c63..a5d2736d3 100644 --- a/src/repository/kodi_repository.rs +++ b/src/repository/kodi_repository.rs @@ -10,10 +10,10 @@ use crate::repository::storage::{ ensure_target_storage_path, get_input_storage_path, hash_bytes, FILE_SUFFIX_DB, }; use crate::repository::xtream_repository::{xtream_get_record_file_path, InputVodInfoRecord}; -use crate::utils::file_lock_manager::FileReadGuard; -use crate::utils::file_utils; -use crate::utils::request_utils::extract_extension_from_url; -use crate::{create_m3u_filter_error_result, info_err, notify_err}; +use crate::utils::file::file_lock_manager::FileReadGuard; +use crate::utils::file::file_utils; +use crate::utils::network::request::extract_extension_from_url; +use crate::m3u_filter_error::{create_m3u_filter_error_result, info_err, notify_err}; use async_std::fs::{create_dir_all, read_dir, remove_dir, remove_file, File}; use async_std::io::{BufReadExt, BufReader, BufWriter, ReadExt, WriteExt}; use chrono::Datelike; @@ -664,7 +664,7 @@ pub async fn kodi_write_strm_playlist( let Some(root_path) = file_utils::get_file_path( &cfg.working_dir, - Some(std::path::PathBuf::from(&output.filename.as_ref().unwrap())), + Some(std::path::PathBuf::from(&output.filename.as_ref().map_or_else(|| "/tmp", |v| v.as_ref()))), ) else { return Err(info_err!(format!( "Failed to get file path for {}", @@ -672,7 +672,7 @@ pub async fn kodi_write_strm_playlist( ))); }; - let credentials_and_server_info = get_credentials_and_server_info(cfg, output); + let credentials_and_server_info = get_credentials_and_server_info(cfg, output).await; let (underscore_whitespace, cleanup, kodi_style) = get_strm_output_options(target); let strm_index_path = strm_get_file_paths(&ensure_target_storage_path(cfg, target.name.as_str())?); @@ -852,19 +852,17 @@ async fn has_strm_file_same_hash(file_path: &PathBuf, content_hash: UUIDType) -> false } -fn get_credentials_and_server_info( +async fn get_credentials_and_server_info( cfg: &Config, output: &TargetOutput, ) -> Option<(ProxyUserCredentials, ApiProxyServerInfo)> { - output - .username - .as_ref() - .and_then(|username| cfg.get_user_credentials(username)) - .filter(|credentials| credentials.proxy == ProxyType::Reverse) - .map(|credentials| { - let server_info = cfg.get_user_server_info(&credentials); - (credentials, server_info) - }) + let username = output.username.as_ref()?; + let credentials = cfg.get_user_credentials(username).await?; + if credentials.proxy != ProxyType::Reverse { + return None; + } + let server_info = cfg.get_user_server_info(&credentials).await; + Some((credentials, server_info)) } async fn read_strm_file_index(strm_file_index_path: &Path) -> std::io::Result> { diff --git a/src/repository/m3u_playlist_iterator.rs b/src/repository/m3u_playlist_iterator.rs index b30d64fcf..d4c24c6c8 100644 --- a/src/repository/m3u_playlist_iterator.rs +++ b/src/repository/m3u_playlist_iterator.rs @@ -1,4 +1,4 @@ -use crate::info_err; +use crate::m3u_filter_error::info_err; use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; use crate::model::api_proxy::{ProxyType, ProxyUserCredentials}; use crate::model::config::{Config, ConfigTarget, ConfigTargetOptions}; @@ -6,7 +6,7 @@ use crate::model::playlist::{M3uPlaylistItem, PlaylistItemType}; use crate::repository::indexed_document::IndexedDocumentIterator; use crate::repository::m3u_repository::m3u_get_file_paths; use crate::repository::storage::ensure_target_storage_path; -use crate::utils::file_lock_manager::FileReadGuard; +use crate::utils::file::file_lock_manager::FileReadGuard; pub const M3U_STREAM_PATH: &str = "m3u-stream"; pub const M3U_RESOURCE_PATH: &str = "resource/m3u"; @@ -46,7 +46,7 @@ impl M3uPlaylistIterator { let include_type_in_url = target_options.is_some_and(|opts| opts.m3u_include_type_in_url); let mask_redirect_url = target_options.is_some_and(|opts| opts.m3u_mask_redirect_url); - let server_info = cfg.get_user_server_info(user); + let server_info = cfg.get_user_server_info(user).await; Ok(Self { reader, base_url: server_info.get_base_url(), diff --git a/src/repository/m3u_repository.rs b/src/repository/m3u_repository.rs index 1329bf40a..7760883fb 100644 --- a/src/repository/m3u_repository.rs +++ b/src/repository/m3u_repository.rs @@ -3,7 +3,7 @@ use std::io::{Error, Write}; use std::path::{Path, PathBuf}; use log::error; -use crate::{create_m3u_filter_error, info_err}; +use crate::m3u_filter_error::{info_err, create_m3u_filter_error}; use crate::m3u_filter_error::{str_to_io_error, M3uFilterError, M3uFilterErrorKind}; use crate::model::api_proxy::{ProxyUserCredentials}; use crate::model::config::{Config, ConfigTarget}; @@ -11,8 +11,8 @@ use crate::model::playlist::{M3uPlaylistItem, PlaylistGroup, PlaylistItem, Playl use crate::repository::indexed_document::{IndexedDocumentDirectAccess, IndexedDocumentWriter}; use crate::repository::m3u_playlist_iterator::M3uPlaylistIterator; use crate::repository::storage::{get_target_storage_path, FILE_SUFFIX_DB, FILE_SUFFIX_INDEX}; -use crate::utils::file_utils; -use crate::utils::file_utils::file_writer; +use crate::utils::file::file_utils; +use crate::utils::file::file_utils::file_writer; const FILE_M3U: &str = "m3u"; macro_rules! cant_write_result { diff --git a/src/repository/playlist_repository.rs b/src/repository/playlist_repository.rs index 5c9de6618..fcb560708 100644 --- a/src/repository/playlist_repository.rs +++ b/src/repository/playlist_repository.rs @@ -1,4 +1,4 @@ -use crate::info_err; +use crate::m3u_filter_error::{info_err}; use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; use crate::model::config::{Config, ConfigTarget, TargetType}; use crate::model::playlist::{PlaylistGroup, PlaylistItemType}; diff --git a/src/repository/storage.rs b/src/repository/storage.rs index 946d36b0a..e5ba9e756 100644 --- a/src/repository/storage.rs +++ b/src/repository/storage.rs @@ -3,8 +3,8 @@ use std::fmt::Write; use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; use crate::model::config::{Config, ConfigInput}; use crate::model::playlist::UUIDType; -use crate::notify_err; -use crate::utils::file_utils; +use crate::m3u_filter_error::{notify_err}; +use crate::utils::file::file_utils; pub(in crate::repository) const FILE_SUFFIX_DB: &str = "db"; pub(in crate::repository) const FILE_SUFFIX_INDEX: &str = "idx"; @@ -44,7 +44,7 @@ pub(in crate::repository) fn get_target_id_mapping_file(target_path: &Path) -> P pub fn ensure_target_storage_path(cfg: &Config, target_name: &str) -> Result { if let Some(path) = get_target_storage_path(cfg, target_name) { if std::fs::create_dir_all(&path).is_err() { - let msg = format!("Failed to save target data, can't create directory {}", &path.to_str().unwrap()); + let msg = format!("Failed to save target data, can't create directory {path:?}"); return Err(notify_err!(msg)); } Ok(path) diff --git a/src/repository/xtream_playlist_iterator.rs b/src/repository/xtream_playlist_iterator.rs index a07a3ce9f..457a2c0da 100644 --- a/src/repository/xtream_playlist_iterator.rs +++ b/src/repository/xtream_playlist_iterator.rs @@ -1,5 +1,5 @@ use log::error; -use crate::info_err; +use crate::m3u_filter_error::info_err; use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; use crate::model::api_proxy::{ProxyUserCredentials}; use crate::model::config::{Config, ConfigTarget}; @@ -7,7 +7,7 @@ use crate::model::playlist::{XtreamCluster, XtreamPlaylistItem}; use crate::model::xtream::XtreamMappingOptions; use crate::repository::indexed_document::{IndexedDocumentIterator}; use crate::repository::xtream_repository::{xtream_get_file_paths, xtream_get_storage_path}; -use crate::utils::file_lock_manager::FileReadGuard; +use crate::utils::file::file_lock_manager::FileReadGuard; pub struct XtreamPlaylistIterator { reader: IndexedDocumentIterator, @@ -35,10 +35,10 @@ impl XtreamPlaylistIterator { .map_err(|err| info_err!(format!("Could not lock document {xtream_path:?}: {err}")))?; let reader = IndexedDocumentIterator::::new(&xtream_path, &idx_path) - .map_err(|err| info_err!(format!("Could not deserialize file {} - {}", &xtream_path.to_str().unwrap(), err)))?; + .map_err(|err| info_err!(format!("Could not deserialize file {xtream_path:?} - {err}")))?; let options = XtreamMappingOptions::from_target_options(target.options.as_ref(), config); - let server_info = config.get_user_server_info(user); + let server_info = config.get_user_server_info(user).await; Ok(Self { reader, options, diff --git a/src/repository/xtream_repository.rs b/src/repository/xtream_repository.rs index 5ec21c058..37aec4673 100644 --- a/src/repository/xtream_repository.rs +++ b/src/repository/xtream_repository.rs @@ -1,5 +1,5 @@ use crate::repository::storage::hex_encode; -use crate::file_utils::file_reader; +use crate::utils::file::file_utils::file_reader; use crate::m3u_filter_error::str_to_io_error; use std::collections::HashMap; use std::fs; @@ -19,9 +19,9 @@ use crate::repository::indexed_document::{IndexedDocumentDirectAccess, IndexedDo use crate::repository::storage::{get_input_storage_path, get_target_id_mapping_file, get_target_storage_path, FILE_SUFFIX_DB, FILE_SUFFIX_INDEX}; use crate::repository::target_id_mapping::{TargetIdMapping, VirtualIdRecord}; use crate::repository::xtream_playlist_iterator::XtreamPlaylistIterator; -use crate::utils::file_utils::open_readonly_file; +use crate::utils::file::file_utils::open_readonly_file; use crate::utils::json_utils::{get_u32_from_serde_value, json_iter_array, json_write_documents_to_file}; -use crate::{create_m3u_filter_error, create_m3u_filter_error_result, info_err, notify_err}; +use crate::m3u_filter_error::{create_m3u_filter_error, create_m3u_filter_error_result, info_err, notify_err}; use crate::utils::hash_utils::generate_playlist_uuid; pub static COL_CAT_LIVE: &str = "cat_live"; @@ -297,7 +297,7 @@ pub async fn xtream_write_playlist( match json_write_documents_to_file(&col_path, data) { Ok(()) => {} Err(err) => { - errors.push(format!("Persisting collection failed: {}: {}", &col_path.to_str().unwrap(), err)); + errors.push(format!("Persisting collection failed: {col_path:?}: {err}")); } } } @@ -577,7 +577,7 @@ pub async fn xtream_load_vod_info( None } -fn rewrite_xtream_vod_info

( +async fn rewrite_xtream_vod_info

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

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

( Ok(result) } -pub fn rewrite_xtream_vod_info_content

( +pub async fn rewrite_xtream_vod_info_content

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

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

( @@ -648,7 +648,7 @@ pub async fn write_and_get_xtream_vod_info

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

( @@ -665,7 +665,7 @@ async fn rewrite_xtream_series_info

( let resource_url = if config.is_reverse_proxy_resource_rewrite_enabled() { match user.proxy { ProxyType::Reverse => { - let server_info = config.get_user_server_info(user); + let server_info = config.get_user_server_info(user).await; let url = server_info.get_base_url(); Some(format!("{url}/resource/series/{}/{}/{}", user.username, user.password, pli.get_virtual_id())) } diff --git a/src/utils/atomic_once_flag.rs b/src/tools/atomic_once_flag.rs similarity index 94% rename from src/utils/atomic_once_flag.rs rename to src/tools/atomic_once_flag.rs index 2f78b5fb7..ba5133f29 100644 --- a/src/utils/atomic_once_flag.rs +++ b/src/tools/atomic_once_flag.rs @@ -20,6 +20,12 @@ pub struct AtomicOnceFlag { ordering: Ordering, } +impl Default for AtomicOnceFlag { + fn default() -> Self { + Self::new() + } +} + impl AtomicOnceFlag { /// Creates a new `AtomicOnceFlag` with the specified memory ordering. pub fn with_ordering(ordering: Ordering) -> Self { diff --git a/src/utils/directed_graph.rs b/src/tools/directed_graph.rs similarity index 97% rename from src/utils/directed_graph.rs rename to src/tools/directed_graph.rs index 17de9006b..3a3a4eb74 100644 --- a/src/utils/directed_graph.rs +++ b/src/tools/directed_graph.rs @@ -199,7 +199,7 @@ where #[cfg(test)] mod tests { - use crate::utils::directed_graph::DirectedGraph; + use crate::tools::directed_graph::DirectedGraph; use std::collections::HashSet; fn are_vecs_equal(vec1: &Vec<&str>, vec2: Vec<&str>) -> bool { @@ -325,4 +325,13 @@ mod tests { assert!(graph.get_dependencies().is_some(), "No dependencies"); } -} \ No newline at end of file +} + +impl Default for DirectedGraph +where + K: Eq + std::hash::Hash + Clone + Display + Debug, +{ + fn default() -> Self { + Self::new() + } +} diff --git a/src/utils/lru_cache.rs b/src/tools/lru_cache.rs similarity index 99% rename from src/utils/lru_cache.rs rename to src/tools/lru_cache.rs index a8615f41f..36b3ba170 100644 --- a/src/utils/lru_cache.rs +++ b/src/tools/lru_cache.rs @@ -1,5 +1,5 @@ use crate::repository::storage::hash_string_as_hex; -use crate::utils::file_utils::traverse_dir; +use crate::utils::file::file_utils::traverse_dir; use crate::utils::size_utils::human_readable_byte_size; use async_std::sync::RwLock; use log::{debug, error, info, trace}; diff --git a/src/tools/mod.rs b/src/tools/mod.rs new file mode 100644 index 000000000..6e3dbf2f0 --- /dev/null +++ b/src/tools/mod.rs @@ -0,0 +1,3 @@ +pub mod directed_graph; +pub mod lru_cache; +pub mod atomic_once_flag; \ No newline at end of file diff --git a/src/utils/compressed_file_reader.rs b/src/utils/compression/compressed_file_reader.rs similarity index 92% rename from src/utils/compressed_file_reader.rs rename to src/utils/compression/compressed_file_reader.rs index 2b71a5957..30f0207f2 100644 --- a/src/utils/compressed_file_reader.rs +++ b/src/utils/compression/compressed_file_reader.rs @@ -1,8 +1,8 @@ use std::io::{BufRead, BufReader, Read, Seek, SeekFrom}; use std::path::Path; use flate2::bufread::{GzDecoder, ZlibDecoder}; -use crate::utils::compression_utils::{is_deflate, is_gzip}; -use crate::utils::file_utils::{file_reader, open_readonly_file}; +use crate::utils::compression::compression_utils::{is_deflate, is_gzip}; +use crate::utils::file::file_utils::{file_reader, open_readonly_file}; pub struct CompressedFileReader { reader: BufReader>, diff --git a/src/utils/compression_utils.rs b/src/utils/compression/compression_utils.rs similarity index 100% rename from src/utils/compression_utils.rs rename to src/utils/compression/compression_utils.rs diff --git a/src/utils/compression/mod.rs b/src/utils/compression/mod.rs new file mode 100644 index 000000000..d28df9813 --- /dev/null +++ b/src/utils/compression/mod.rs @@ -0,0 +1,2 @@ +pub mod compressed_file_reader; +pub mod compression_utils; diff --git a/src/utils/config_reader.rs b/src/utils/config_reader.rs index 190021d02..2f8a55a82 100644 --- a/src/utils/config_reader.rs +++ b/src/utils/config_reader.rs @@ -6,12 +6,12 @@ use chrono::Local; use log::{debug, error, info, warn}; use regex::Regex; use serde::Serialize; -use crate::{create_m3u_filter_error, create_m3u_filter_error_result, exit, handle_m3u_filter_error_result, info_err}; -use crate::m3u_filter_error::{to_io_error, M3uFilterError, M3uFilterErrorKind}; +use crate::utils::sys_utils::exit; +use crate::m3u_filter_error::{to_io_error, M3uFilterError, M3uFilterErrorKind, create_m3u_filter_error, create_m3u_filter_error_result, info_err, handle_m3u_filter_error_result}; use crate::model::api_proxy::ApiProxyConfig; use crate::model::config::{Config, ConfigDto}; use crate::model::mapping::Mappings; -use crate::utils::{file_utils, multi_file_reader}; +use crate::utils::file::{file_utils, multi_file_reader}; pub fn read_mappings(args_mapping: Option, cfg: &mut Config) -> Result, M3uFilterError> { let mappings_file: String = args_mapping.unwrap_or_else(|| file_utils::get_default_mappings_path(cfg.t_config_path.as_str())); diff --git a/src/utils/download.rs b/src/utils/download.rs deleted file mode 100644 index 45d0831ff..000000000 --- a/src/utils/download.rs +++ /dev/null @@ -1,18 +0,0 @@ -use std::borrow::Cow; -use crate::utils::{file_utils}; -use std::path::PathBuf; -use crate::debug_if_enabled; - - -pub fn prepare_file_path(persist: Option<&str>, working_dir: &str, action: &str) -> Option { - let persist_file: Option = - persist.map(|persist_path| file_utils::prepare_persist_path(persist_path, action)); - if persist_file.is_some() { - let file_path = file_utils::get_file_path(working_dir, persist_file); - debug_if_enabled!("persist to file: {}", file_path.as_ref().map_or(Cow::from("?"), |p| p.to_string_lossy())); - file_path - } else { - None - } -} - diff --git a/src/utils/file_lock_manager.rs b/src/utils/file/file_lock_manager.rs similarity index 100% rename from src/utils/file_lock_manager.rs rename to src/utils/file/file_lock_manager.rs diff --git a/src/utils/file_utils.rs b/src/utils/file/file_utils.rs similarity index 92% rename from src/utils/file_utils.rs rename to src/utils/file/file_utils.rs index 99b74486e..da6bb85c3 100644 --- a/src/utils/file_utils.rs +++ b/src/utils/file/file_utils.rs @@ -1,3 +1,4 @@ +use std::borrow::Cow; use std::fs; use std::fs::{File, OpenOptions}; use std::io::{BufReader, BufWriter, Read, Write}; @@ -5,6 +6,7 @@ use std::path::{Path, PathBuf}; use log::{debug, error}; use path_clean::PathClean; +use crate::utils::debug_if_enabled; use crate::m3u_filter_error::str_to_io_error; const USER_FILE: &str = "user.txt"; @@ -14,14 +16,6 @@ const SOURCE_FILE: &str = "source.yml"; const MAPPING_FILE: &str = "mapping.yml"; const API_PROXY_FILE: &str = "api-proxy.yml"; -#[macro_export] -macro_rules! exit { - ($($arg:tt)*) => {{ - error!($($arg)*); - std::process::exit(1); - }}; -} - pub fn file_writer(w: W) -> BufWriter where W: Write { @@ -244,4 +238,16 @@ where } Ok(()) +} + +pub fn prepare_file_path(persist: Option<&str>, working_dir: &str, action: &str) -> Option { + let persist_file: Option = + persist.map(|persist_path| prepare_persist_path(persist_path, action)); + if persist_file.is_some() { + let file_path = get_file_path(working_dir, persist_file); + debug_if_enabled!("persist to file: {}", file_path.as_ref().map_or(Cow::from("?"), |p| p.to_string_lossy())); + file_path + } else { + None + } } \ No newline at end of file diff --git a/src/utils/file/mod.rs b/src/utils/file/mod.rs new file mode 100644 index 000000000..8afe663ac --- /dev/null +++ b/src/utils/file/mod.rs @@ -0,0 +1,3 @@ +pub mod file_utils; +pub mod multi_file_reader; +pub mod file_lock_manager; \ No newline at end of file diff --git a/src/utils/multi_file_reader.rs b/src/utils/file/multi_file_reader.rs similarity index 97% rename from src/utils/multi_file_reader.rs rename to src/utils/file/multi_file_reader.rs index d595fe740..da4a0bdac 100644 --- a/src/utils/multi_file_reader.rs +++ b/src/utils/file/multi_file_reader.rs @@ -1,7 +1,7 @@ use std::fs::File; use std::io::{self, BufReader, ErrorKind, Read}; use std::path::{PathBuf}; -use crate::utils::file_utils::file_reader; +use crate::utils::file::file_utils::file_reader; pub struct MultiFileReader { files: Vec, diff --git a/src/utils/file_reader.rs b/src/utils/file_reader.rs deleted file mode 100644 index f45861b8e..000000000 --- a/src/utils/file_reader.rs +++ /dev/null @@ -1,24 +0,0 @@ -use std::fs::File; -use linereader::LineReader; - -pub struct FileReader { - reader: LineReader, -} - -impl FileReader { - pub fn new(file: File) -> Self { - Self { - reader: LineReader::new(file), - } - } -} - -impl Iterator for FileReader { - type Item = String; - fn next(&mut self) -> Option { - if let Some(Ok(buf)) = self.reader.next_line() { - return Some(String::from_utf8_lossy(buf).trim_end_matches(char::is_control).to_string()); - } - None - } -} \ No newline at end of file diff --git a/src/utils/json_utils.rs b/src/utils/json_utils.rs index c8e26e2f3..b01968b80 100644 --- a/src/utils/json_utils.rs +++ b/src/utils/json_utils.rs @@ -6,7 +6,7 @@ use std::path::Path; use serde::de::DeserializeOwned; use serde::{Deserialize, Serialize}; use serde_json::{self, Deserializer, Value}; -use crate::utils::file_utils::{file_reader, file_writer}; +use crate::utils::file::file_utils::{file_reader, file_writer}; fn read_skipping_ws(mut reader: impl Read) -> io::Result { loop { @@ -64,7 +64,7 @@ pub fn json_iter_array( std::iter::from_fn(move || yield_next_obj(&mut reader, &mut at_start).transpose()) } -pub fn json_filter_file(file_path: &Path, filter: &HashMap<&str, &str>) -> Vec { +pub fn json_filter_file(file_path: &Path, filter: &HashMap<&str, &str, S>) -> Vec { let mut filtered: Vec = Vec::with_capacity(1024); if !file_path.exists() { return filtered; // Return early if the file does not exist diff --git a/src/utils/mod.rs b/src/utils/mod.rs index b5e204c9d..5e8613121 100644 --- a/src/utils/mod.rs +++ b/src/utils/mod.rs @@ -1,23 +1,13 @@ -pub mod file_utils; -pub mod request_utils; -pub mod download; pub mod string_utils; pub mod json_utils; pub mod config_reader; pub mod default_utils; -pub mod multi_file_reader; -pub mod file_lock_manager; -pub mod compressed_file_reader; -mod compression_utils; -pub mod directed_graph; -pub mod lru_cache; pub mod size_utils; -pub mod sys; -pub mod atomic_once_flag; -pub mod xtream_utils; -pub mod m3u_utils; -pub mod epg_utils; +pub mod sys_utils; pub mod hash_utils; +pub mod compression; +pub(crate) mod file; +pub(crate) mod network; #[macro_export] macro_rules! debug_if_enabled { @@ -48,3 +38,6 @@ macro_rules! trace_if_enabled { } }; } + +pub use debug_if_enabled; +pub use trace_if_enabled; \ No newline at end of file diff --git a/src/utils/epg_utils.rs b/src/utils/network/epg.rs similarity index 77% rename from src/utils/epg_utils.rs rename to src/utils/network/epg.rs index 2ce4ae301..cbea69a25 100644 --- a/src/utils/epg_utils.rs +++ b/src/utils/network/epg.rs @@ -3,8 +3,9 @@ use log::debug; use crate::m3u_filter_error::M3uFilterError; use crate::model::config::{Config, ConfigInput}; use crate::model::xmltv::TVGuide; -use crate::utils::download::prepare_file_path; -use crate::utils::{file_utils, request_utils}; +use crate::utils::file::file_utils::prepare_file_path; +use crate::utils::network::request; +use crate::utils::file::file_utils; pub async fn get_xmltv(client: Arc, _cfg: &Config, input: &ConfigInput, working_dir: &str) -> (Option, Vec) { match &input.epg_url { @@ -14,7 +15,7 @@ pub async fn get_xmltv(client: Arc, _cfg: &Config, input: &Conf let persist_file_path = prepare_file_path(input.persist.as_deref(), working_dir, "") .map(|path| file_utils::add_prefix_to_filename(&path, "epg_", Some("xml"))); - match request_utils::get_input_text_content_as_file(client, input, working_dir, url, persist_file_path).await { + match request::get_input_text_content_as_file(client, input, working_dir, url, persist_file_path).await { Ok(file) => { (Some(TVGuide { file }), vec![]) } diff --git a/src/utils/m3u_utils.rs b/src/utils/network/m3u.rs similarity index 63% rename from src/utils/m3u_utils.rs rename to src/utils/network/m3u.rs index 92ecf0d60..06fe9cd0d 100644 --- a/src/utils/m3u_utils.rs +++ b/src/utils/network/m3u.rs @@ -2,16 +2,16 @@ use std::sync::Arc; use crate::m3u_filter_error::M3uFilterError; use crate::model::config::{Config, ConfigInput}; use crate::model::playlist::PlaylistGroup; -use crate::processing::m3u_parser; -use crate::utils::download::prepare_file_path; -use crate::utils::request_utils; +use crate::processing::parser::m3u; +use crate::utils::file::file_utils::prepare_file_path; +use crate::utils::network::request; pub async fn get_m3u_playlist(client: Arc, cfg: &Config, input: &ConfigInput, working_dir: &str) -> (Vec, Vec) { let url = input.url.clone(); let persist_file_path = prepare_file_path(input.persist.as_deref(), working_dir, ""); - match request_utils::get_input_text_content(client, input, working_dir, &url, persist_file_path).await { + match request::get_input_text_content(client, input, working_dir, &url, persist_file_path).await { Ok(text) => { - (m3u_parser::parse_m3u(cfg, input, text.lines()), vec![]) + (m3u::parse_m3u(cfg, input, text.lines()), vec![]) } Err(err) => (vec![], vec![err]) } diff --git a/src/utils/network/mod.rs b/src/utils/network/mod.rs new file mode 100644 index 000000000..ce36c706e --- /dev/null +++ b/src/utils/network/mod.rs @@ -0,0 +1,4 @@ +pub mod request; +pub mod xtream; +pub mod m3u; +pub mod epg; \ No newline at end of file diff --git a/src/utils/request_utils.rs b/src/utils/network/request.rs similarity index 98% rename from src/utils/request_utils.rs rename to src/utils/network/request.rs index c5bbb1a01..f23206c24 100644 --- a/src/utils/request_utils.rs +++ b/src/utils/network/request.rs @@ -21,9 +21,10 @@ use crate::model::config::ConfigInput; use crate::model::stats::format_elapsed_time; use crate::repository::storage::get_input_storage_path; use crate::repository::xtream_repository::FILE_EPG; -use crate::utils::compression_utils::{is_deflate, is_gzip, ENCODING_DEFLATE, ENCODING_GZIP}; -use crate::utils::file_utils::{get_file_path, persist_file}; -use crate::{create_m3u_filter_error_result, debug_if_enabled}; +use crate::utils::compression::compression_utils::{is_deflate, is_gzip, ENCODING_DEFLATE, ENCODING_GZIP}; +use crate::utils::file::file_utils::{get_file_path, persist_file}; +use crate::m3u_filter_error::create_m3u_filter_error_result; +use crate::utils::debug_if_enabled; pub const fn bytes_to_megabytes(bytes: u64) -> u64 { bytes / 1_048_576 @@ -440,7 +441,7 @@ pub fn replace_extension(path: &str, new_ext: &str) -> String { #[cfg(test)] mod tests { - use crate::utils::request_utils::{replace_extension, sanitize_sensitive_info}; + use crate::utils::network::request::{replace_extension, sanitize_sensitive_info}; #[test] fn test_url_mask() { diff --git a/src/utils/xtream_utils.rs b/src/utils/network/xtream.rs similarity index 88% rename from src/utils/xtream_utils.rs rename to src/utils/network/xtream.rs index ecb25337b..29ed33bec 100644 --- a/src/utils/xtream_utils.rs +++ b/src/utils/network/xtream.rs @@ -2,16 +2,16 @@ use crate::Arc; use crate::m3u_filter_error::{str_to_io_error, M3uFilterError}; use crate::model::config::{Config, ConfigInput, ConfigTarget}; use crate::model::playlist::{PlaylistEntry, PlaylistGroup, XtreamCluster, XtreamPlaylistItem}; -use crate::processing::{xtream_parser}; +use crate::processing::parser::xtream; use crate::repository::xtream_repository::{rewrite_xtream_series_info_content, rewrite_xtream_vod_info_content, xtream_get_input_info}; use crate::repository::xtream_repository; -use crate::utils::{request_utils}; use log::{info, warn}; use std::cmp::Ordering; -use std::io::{Error}; +use std::io::Error; use crate::model::api_proxy::{ProxyType, ProxyUserCredentials}; use crate::utils::json_utils::get_string_from_serde_value; -use crate::utils::request_utils::extract_extension_from_url; +use crate::utils::network::request; +use crate::utils::network::request::extract_extension_from_url; pub const ACTION_GET_SERIES_INFO: &str = "get_series_info"; pub const ACTION_GET_VOD_INFO: &str = "get_vod_info"; @@ -53,7 +53,7 @@ pub fn get_xtream_player_api_info_url(input: &ConfigInput, cluster: XtreamCluste pub async fn get_xtream_stream_info_content(client: Arc, info_url: &str, input: &ConfigInput) -> Result { - request_utils::download_text_content(client, input, info_url, None).await + request::download_text_content(client, input, info_url, None).await } #[allow(clippy::too_many_arguments)] @@ -86,7 +86,7 @@ where } else if cluster == XtreamCluster::Video { if let Some(content) = xtream_repository::xtream_load_vod_info(config, target.name.as_str(), pli.get_virtual_id()).await { // Deliver existing target content - return rewrite_xtream_vod_info_content(config, target, pli, user, &content); + return rewrite_xtream_vod_info_content(config, target, pli, user, &content).await; } // Check if the content has been resolved let resolve_vod = target.options.as_ref().is_some_and(|opt| opt.xtream_resolve_vod); @@ -142,7 +142,7 @@ pub async fn get_xtream_playlist(client: Arc, input: &ConfigInp let base_url = get_xtream_stream_url_base(&input.url, username, password); - if let Err(err) = request_utils::get_input_json_content(Arc::clone(&client), input, base_url.as_str(), None).await { + if let Err(err) = request::get_input_json_content(Arc::clone(&client), input, base_url.as_str(), None).await { warn!("Failed to login xtream account {username} {err}"); return (Vec::with_capacity(0), vec![err]); }; @@ -156,18 +156,18 @@ pub async fn get_xtream_playlist(client: Arc, input: &ConfigInp if !skip_cluster.contains(xtream_cluster) { let category_url = format!("{base_url}&action={category}"); let stream_url = format!("{base_url}&action={stream}"); - let category_file_path = crate::utils::download::prepare_file_path(input.persist.as_deref(), working_dir, format!("{category}_").as_str()); - let stream_file_path = crate::utils::download::prepare_file_path(input.persist.as_deref(), working_dir, format!("{stream}_").as_str()); + let category_file_path = crate::utils::file::file_utils::prepare_file_path(input.persist.as_deref(), working_dir, format!("{category}_").as_str()); + let stream_file_path = crate::utils::file::file_utils::prepare_file_path(input.persist.as_deref(), working_dir, format!("{stream}_").as_str()); match futures::join!( - request_utils::get_input_json_content(Arc::clone(&client), input, category_url.as_str(), category_file_path), - request_utils::get_input_json_content(Arc::clone(&client), input, stream_url.as_str(), stream_file_path) + request::get_input_json_content(Arc::clone(&client), input, category_url.as_str(), category_file_path), + request::get_input_json_content(Arc::clone(&client), input, stream_url.as_str(), stream_file_path) ) { (Ok(category_content), Ok(stream_content)) => { - match xtream_parser::parse_xtream(input, - *xtream_cluster, - &category_content, - &stream_content) { + match xtream::parse_xtream(input, + *xtream_cluster, + &category_content, + &stream_content) { Ok(sub_playlist_parsed) => { if let Some(mut xtream_sub_playlist) = sub_playlist_parsed { playlist_groups.append(&mut xtream_sub_playlist); diff --git a/src/utils/sys.rs b/src/utils/sys_utils.rs similarity index 94% rename from src/utils/sys.rs rename to src/utils/sys_utils.rs index 2bf8f1fa0..1a93b92bd 100644 --- a/src/utils/sys.rs +++ b/src/utils/sys_utils.rs @@ -1,3 +1,12 @@ +#[macro_export] +macro_rules! exit { + ($($arg:tt)*) => {{ + error!($($arg)*); + std::process::exit(1); + }}; +} +pub use exit; + #[cfg(target_os = "linux")] fn get_memory_usage_linux() -> std::io::Result { use std::fs::File;