diff --git a/README.md b/README.md index 246cb97bb..eaf80548b 100644 --- a/README.md +++ b/README.md @@ -612,6 +612,15 @@ sort: order: asc channels: - { field: name, group_pattern: '^DE.*', order: asc } + - field: name + group_pattern: '^FR.*' + order: asc + sequence: + - TF1 + - TF1+1 + - FRANCE 2 + - FRANCE 3 + - FRANCE 4 ``` ### 2.2.2.2 `output` @@ -1360,7 +1369,7 @@ Image docker build -t m3u-filter . ``` docker-compose.yml -```dockerfile +```docker version: '3' services: m3u-filter: @@ -1383,7 +1392,7 @@ This example is for the local image, the official can be found under `ghcr.io/eu If you want to use m3u-filter with docker-compose, there is a `--healthcheck` argument for healthchecks -```dockerfile +```docker healthcheck: test: ["CMD", "/m3u-filter", "-p", "/config" "--healthcheck"] interval: 30s diff --git a/src/api/api_utils.rs b/src/api/api_utils.rs index 70f368381..3cfa59488 100644 --- a/src/api/api_utils.rs +++ b/src/api/api_utils.rs @@ -10,13 +10,13 @@ use crate::api::model::streams::provider_stream::{create_custom_video_stream_res use crate::api::model::streams::provider_stream_factory::BufferStreamOptions; use crate::api::model::streams::shared_stream_manager::SharedStreamManager; use crate::auth::authenticator::Claims; -use crate::model::api_proxy::{ProxyUserCredentials, UserConnectionPermission}; -use crate::model::config::{ConfigInput, ConfigTarget, InputFetchMethod}; -use crate::model::playlist::PlaylistItemType; +use crate::model::api_proxy::{ProxyType, ProxyUserCredentials, UserConnectionPermission}; +use crate::model::config::{ConfigInput, ConfigTarget, InputFetchMethod, TargetType}; +use crate::model::playlist::{PlaylistEntry, PlaylistItemType, XtreamCluster}; use crate::tools::lru_cache::LRUResourceCache; use crate::utils::file::file_utils::create_new_file_for_write; use crate::utils::network::request; -use crate::utils::network::request::sanitize_sensitive_info; +use crate::utils::network::request::{extract_extension_from_url, replace_url_extension, sanitize_sensitive_info}; use crate::utils::{debug_if_enabled, trace_if_enabled}; use axum::http::HeaderMap; use axum::response::IntoResponse; @@ -71,9 +71,11 @@ macro_rules! try_result_bad_request { pub use try_option_bad_request; pub use try_result_bad_request; +use crate::api::endpoints::xtream_api::{get_xtream_player_api_stream_url, XtreamApiStreamContext}; use crate::api::model::stream::{BoxedProviderStream, ProviderStreamInfo, ProviderStreamResponse}; use crate::api::model::streams::throttled_stream::ThrottledStream; use crate::tools::atomic_once_flag::AtomicOnceFlag; +use crate::utils::constants::{DASH_EXT, HLS_EXT}; use crate::utils::default_utils::default_grace_period_millis; #[allow(clippy::missing_panics_doc)] @@ -314,6 +316,102 @@ async fn create_stream_response_details(app_state: &AppState, stream_options: &S } } +pub struct RedirectParams<'a, P> +where + P: PlaylistEntry +{ + pub item: &'a P, + pub provider_id: Option, + pub cluster: XtreamCluster, + pub target_type: TargetType, + pub target: &'a ConfigTarget, + pub input: Option<&'a ConfigInput>, + pub user: &'a ProxyUserCredentials, + pub stream_ext: Option<&'a str>, + pub req_context: XtreamApiStreamContext, + pub action_path: &'a str, +} + +pub fn redirect_response

(params: &RedirectParams

) -> Option +where + P: PlaylistEntry +{ + + let item_type = params.item.get_item_type(); + let provider_url = ¶ms.item.get_provider_url(); + + let redirect_request = params.user.proxy == ProxyType::Redirect || params.target.is_force_redirect(item_type); + let is_hls_request = item_type == PlaylistItemType::LiveHls || params.stream_ext == Some(HLS_EXT); + let is_dash_request = !is_hls_request && item_type == PlaylistItemType::LiveDash || params.stream_ext == Some(DASH_EXT); + + if params.target_type == TargetType::M3u { + if redirect_request || is_dash_request { + let redirect_url = if is_hls_request { &replace_url_extension(provider_url, HLS_EXT) } else { provider_url }; + let redirect_url = if is_dash_request { &replace_url_extension(redirect_url, DASH_EXT) } else { redirect_url }; + // TODO alias processing, redirect to different aliases + debug_if_enabled!("Redirecting stream request to {}", sanitize_sensitive_info(redirect_url)); + return Some(redirect(redirect_url.as_str()).into_response()); + } + } else if params.target_type == TargetType::Xtream { + + let Some(provider_id) = params.provider_id else { + return Some(StatusCode::BAD_REQUEST.into_response()); + }; + + let Some(input) = params.input else { + return Some(StatusCode::BAD_REQUEST.into_response()); + }; + + // handle redirect for series + if redirect_request && params.cluster == XtreamCluster::Series { + let ext = params.stream_ext.unwrap_or_default(); + let url = input.url.as_str(); + let username = input.username.as_ref().map_or("", |v| v); + let password = input.password.as_ref().map_or("", |v| v); + // TODO do i need action_path like for timeshift ? + let stream_url = format!("{url}/series/{username}/{password}/{provider_id}{ext}"); + return Some(redirect(&stream_url).into_response()); + } + + + let extension = params.stream_ext.map_or_else( + || extract_extension_from_url(provider_url).map_or_else(String::new, std::string::ToString::to_string), + std::string::ToString::to_string); + + // if there is a action_path (like for timeshift duration/start) it will be added in front of the stream_id + let query_path = if params.action_path.is_empty() { + format!("{provider_id}{extension}") + } else { + format!("{}/{provider_id}{extension}", params.action_path) + }; + + let target_name = params.target.name.as_str(); + let virtual_id = params.item.get_virtual_id(); + let stream_url = match get_xtream_player_api_stream_url(input, ¶ms.req_context, &query_path, provider_url) { + None => { + error!("Cant find stream url for target {target_name}, context {}, stream_id {virtual_id}", params.req_context); + return Some(StatusCode::BAD_REQUEST.into_response()) + } + Some(url) => url, + }; + + // hls or dash redirect + if redirect_request || is_dash_request { + let redirect_url = if is_hls_request { &replace_url_extension(&stream_url, HLS_EXT) } else { &replace_url_extension(&stream_url, DASH_EXT) }; + debug_if_enabled!("Redirecting stream request to {}", sanitize_sensitive_info(redirect_url)); + return Some(redirect(redirect_url).into_response()); + } + + if redirect_request { + debug_if_enabled!("Redirecting stream request to {}", sanitize_sensitive_info(&stream_url)); + return Some(redirect(&stream_url).into_response()); + } + + } + + None +} + /// # Panics #[allow(clippy::too_many_arguments)] pub async fn stream_response(app_state: &AppState, diff --git a/src/api/endpoints/m3u_api.rs b/src/api/endpoints/m3u_api.rs index fd137a8d6..3495a6d7d 100644 --- a/src/api/endpoints/m3u_api.rs +++ b/src/api/endpoints/m3u_api.rs @@ -1,21 +1,21 @@ -use crate::api::api_utils::{get_user_target, get_user_target_by_credentials, redirect, resource_response, separate_number_and_remainder, stream_response, try_option_bad_request, try_result_bad_request}; +use crate::api::api_utils::{get_user_target, get_user_target_by_credentials, redirect, redirect_response, resource_response, separate_number_and_remainder, stream_response, try_option_bad_request, try_result_bad_request, RedirectParams}; 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, UserConnectionPermission}; use crate::model::config::{TargetType}; -use crate::model::playlist::{FieldGetAccessor, PlaylistItemType, XtreamCluster}; +use crate::model::playlist::{FieldGetAccessor, PlaylistEntry, PlaylistItemType, XtreamCluster}; use crate::repository::m3u_repository::{m3u_get_item_for_stream_id, m3u_load_rewrite_playlist}; -use crate::utils::network::request::{replace_url_extension, sanitize_sensitive_info}; -use crate::utils::debug_if_enabled; +use crate::utils::network::request::{sanitize_sensitive_info}; use axum::response::IntoResponse; use bytes::Bytes; use futures::stream; use log::{debug, error}; use std::sync::Arc; +use crate::api::endpoints::xtream_api::XtreamApiStreamContext; use crate::api::model::streams::provider_stream::{create_custom_video_stream_response, CustomVideoStreamType}; use crate::repository::storage_const; -use crate::utils::constants::{DASH_EXT, HLS_EXT}; +use crate::utils::constants::{HLS_EXT}; async fn m3u_api( api_req: &UserApiRequest, @@ -92,17 +92,27 @@ async fn m3u_api_stream( let input = app_state.config.get_input_by_name(m3u_item.input_name.as_str()); - let is_hls_request = m3u_item.item_type == PlaylistItemType::LiveHls || stream_ext.as_deref() == Some(HLS_EXT); - let is_dash_request = !is_hls_request && m3u_item.item_type == PlaylistItemType::LiveDash || stream_ext.as_deref() == Some(DASH_EXT); + let cluster = XtreamCluster::try_from(m3u_item.item_type).unwrap_or(XtreamCluster::Live); + let context = XtreamApiStreamContext::try_from(cluster).unwrap_or(XtreamApiStreamContext::Live); - if user.proxy == ProxyType::Redirect || is_dash_request || target.is_force_redirect(m3u_item.item_type) { - let redirect_url = if is_hls_request { &replace_url_extension(&m3u_item.url, HLS_EXT) } else { &m3u_item.url }; - let redirect_url = if is_dash_request { &replace_url_extension(redirect_url, DASH_EXT) } else { redirect_url }; - // TODO alias processing, redirect to different aliases - debug_if_enabled!("Redirecting m3u stream request to {}", sanitize_sensitive_info(redirect_url)); - return redirect(redirect_url.as_str()).into_response(); + let redirect_params = RedirectParams { + item: &m3u_item, + provider_id: m3u_item.get_provider_id(), + cluster, + target_type: TargetType::Xtream, + target, + input, + user: &user, + stream_ext: stream_ext.as_deref(), + req_context: context, + action_path: "" // TODO is there timeshoft or something like that ? + }; + + if let Some(response) = redirect_response(&redirect_params) { + return response.into_response(); } + let is_hls_request = m3u_item.item_type == PlaylistItemType::LiveHls || stream_ext.as_deref() == Some(HLS_EXT); // Reverse proxy mode if is_hls_request { let target_name = &target.name; diff --git a/src/api/endpoints/xtream_api.rs b/src/api/endpoints/xtream_api.rs index 1132e79a1..8ac0bbb2f 100644 --- a/src/api/endpoints/xtream_api.rs +++ b/src/api/endpoints/xtream_api.rs @@ -1,7 +1,7 @@ // https://github.com/tellytv/go.xtream-codes/blob/master/structs.go use crate::api::api_utils; -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::api_utils::{get_user_target, get_user_target_by_credentials, redirect_response, resource_response, separate_number_and_remainder, serve_file, stream_response, RedirectParams}; use crate::api::api_utils::{redirect, try_option_bad_request, try_result_bad_request}; use crate::api::endpoints::hls_api::handle_hls_stream_request; use crate::api::endpoints::xmltv_api::get_empty_epg_response; @@ -20,12 +20,12 @@ use crate::model::xtream_const; use crate::repository::playlist_repository::get_target_id_mapping; use crate::repository::storage::{get_target_storage_path, hex_encode}; use crate::repository::{storage_const, user_repository, xtream_repository}; -use crate::utils::constants::{DASH_EXT, HLS_EXT}; +use crate::utils::constants::{HLS_EXT}; use crate::utils::debug_if_enabled; use crate::utils::hash_utils::generate_playlist_uuid; use crate::utils::json_utils; use crate::utils::json_utils::get_u32_from_serde_value; -use crate::utils::network::request::{extract_extension_from_url, replace_url_extension, sanitize_sensitive_info}; +use crate::utils::network::request::{extract_extension_from_url, sanitize_sensitive_info}; use crate::utils::network::xtream::create_vod_info_from_item; use crate::utils::network::{request, xtream}; use crate::utils::trace_if_enabled; @@ -44,7 +44,7 @@ use std::str::FromStr; use std::sync::Arc; #[derive(Serialize, Deserialize, Debug, Clone, Eq, PartialEq)] -pub(in crate::api) enum XtreamApiStreamContext { +pub enum XtreamApiStreamContext { LiveAlt, Live, Movie, @@ -197,50 +197,41 @@ async fn xtream_player_api_stream( let (action_stream_id, stream_ext) = separate_number_and_remainder(stream_req.stream_id); let virtual_id: u32 = try_result_bad_request!(action_stream_id.trim().parse()); let (pli, mapping) = try_result_bad_request!(xtream_repository::xtream_get_item_for_stream_id(virtual_id, &app_state.config, target, None), true, format!("Failed to read xtream item for stream id {}", virtual_id)); + let input = try_option_bad_request!(app_state.config.get_input_by_name(pli.input_name.as_str()), true, format!("Cant find input for target {target_name}, context {}, stream_id {virtual_id}", stream_req.context)); - - let redirect_request = user.proxy == ProxyType::Redirect || target.is_force_redirect(pli.item_type); - - // handle redirect for series - if redirect_request && pli.xtream_cluster == XtreamCluster::Series { - let ext = stream_ext.unwrap_or_else(String::new); - let url = input.url.as_str(); - let username = input.username.as_ref().map_or("", |v| v); - let password = input.password.as_ref().map_or("", |v| v); - // TODO do i need action_path like for timeshift ? - let stream_url = format!("{url}/series/{username}/{password}/{}{ext}", mapping.provider_id); - return redirect(&stream_url).into_response(); + let redirect_params = RedirectParams { + item: &pli, + provider_id: Some(mapping.provider_id), + cluster: pli.xtream_cluster, + target_type: TargetType::Xtream, + target, + input: Some(input), + user: &user, + stream_ext: stream_ext.as_deref(), + req_context: stream_req.context.clone(), + action_path: stream_req.action_path + }; + if let Some(response) = redirect_response(&redirect_params) { + return response.into_response(); } let extension = stream_ext.unwrap_or_else( || extract_extension_from_url(&pli.url).map_or_else(String::new, std::string::ToString::to_string)); - // if there is a action_path (like for timeshift duration/start) it will be added in front of the stream_id let query_path = if stream_req.action_path.is_empty() { format!("{}{extension}", pli.provider_id) } else { format!("{}/{}{extension}", stream_req.action_path, pli.provider_id) }; - + let stream_url = try_option_bad_request!(get_xtream_player_api_stream_url(input, - &stream_req.context, &query_path, pli.url.as_str()), + &stream_req.context, &query_path, &pli.url), true, format!("Cant find stream url for target {target_name}, context {}, stream_id {virtual_id}", stream_req.context)); + let is_hls_request = pli.item_type == PlaylistItemType::LiveHls || extension == HLS_EXT; - // hls or dash redirect - if redirect_request && (is_hls_request || pli.item_type == PlaylistItemType::LiveDash) { - let redirect_url = if is_hls_request { &replace_url_extension(&stream_url, HLS_EXT) } else { &replace_url_extension(&stream_url, DASH_EXT) }; - debug_if_enabled!("Redirecting stream request to {}", sanitize_sensitive_info(redirect_url)); - return redirect(redirect_url).into_response(); - } - - if redirect_request { - debug_if_enabled!("Redirecting stream request to {}", sanitize_sensitive_info(&stream_url)); - return redirect(&stream_url).into_response(); - } - // Reverse proxy mode if is_hls_request { return handle_hls_stream_request(app_state, &user, &stream_url, pli.virtual_id, input).await.into_response(); diff --git a/src/api/model/active_provider_manager.rs b/src/api/model/active_provider_manager.rs index 0d64533ff..383a36e43 100644 --- a/src/api/model/active_provider_manager.rs +++ b/src/api/model/active_provider_manager.rs @@ -3,6 +3,7 @@ use std::collections::HashMap; use std::ops::Deref; use std::sync::{Arc}; use std::sync::atomic::{AtomicU16, AtomicUsize, Ordering}; +use log::{debug, log_enabled}; use tokio::sync::RwLock; pub struct ProviderConnectionGuard { @@ -122,6 +123,18 @@ impl ProviderConfig { 3 } + // is intended to use with redirects, to cycle through provider + fn get_next(&self, grace: bool) -> bool { + let connections = self.current_connections.load(Ordering::SeqCst); + if self.max_connections == 0 { + return true; + } + if (!grace && connections < self.max_connections) || (grace && connections <= self.max_connections) { + return true; + } + false + } + pub fn release(&self) { let connections = self.current_connections.load(Ordering::SeqCst); if connections > 0 { @@ -154,6 +167,13 @@ impl ProviderConfigWrapper { _ => ProviderAllocation::Exhausted, } } + + pub fn get_next(&self, grace: bool) -> Option> { + if self.inner.get_next(grace) { + return Some(Arc::clone(&self.inner)); + } + None + } } impl Deref for ProviderConfigWrapper { type Target = ProviderConfig; @@ -174,6 +194,14 @@ enum ProviderLineup { } impl ProviderLineup { + + fn get_next(&self) -> Option> { + match self { + ProviderLineup::Single(lineup) => lineup.get_next(), + ProviderLineup::Multi(lineup) => lineup.get_next(), + } + } + fn acquire(&self) -> ProviderAllocation { match self { ProviderLineup::Single(lineup) => lineup.acquire(), @@ -202,6 +230,10 @@ impl SingleProviderLineup { } } + fn get_next(&self) -> Option> { + self.provider.get_next(true) + } + fn acquire(&self) -> ProviderAllocation { self.provider.try_allocate(true) } @@ -338,6 +370,31 @@ impl MultiProviderLineup { ProviderAllocation::Exhausted } + // Used for redirect to cylce through provider + fn get_next_provider_from_group(priority_group: &ProviderPriorityGroup, grace: bool) -> Option> { + match priority_group { + ProviderPriorityGroup::SingleProviderGroup(p) => { + return p.get_next(grace); + } + ProviderPriorityGroup::MultiProviderGroup(index, pg) => { + let mut idx = index.load(Ordering::SeqCst); + let provider_count = pg.len(); + let start = idx; + for _ in start..provider_count { + let p = pg.get(idx).unwrap(); + idx = (idx + 1) % provider_count; + let result = p.get_next(grace); + if result.is_some() { + index.store(idx, Ordering::SeqCst); + return result; + } + } + index.store(idx, Ordering::SeqCst); + } + } + None + } + /// Attempts to acquire a provider from the lineup based on priority and availability. /// /// # Returns @@ -390,20 +447,35 @@ impl MultiProviderLineup { } ProviderAllocation::Exhausted - // let provider = &self.providers[main_idx]; - // self.index.store((main_idx + 1) % provider_count, Ordering::SeqCst); - // - // match provider { - // ProviderPriorityGroup::SingleProviderGroup(p) => ProviderAllocation::Available(p), - // ProviderPriorityGroup::MultiProviderGroup(gindex, group) => { - // let idx = gindex.load(Ordering::SeqCst); - // gindex.store((idx + 1) % group.len(), Ordering::SeqCst); - // match group.get(idx) { - // None => ProviderAllocation::Exhausted, - // Some(p) => ProviderAllocation::Available(p) - // } - // } - // } + } + + // it intended to use with redirects to cycle through provider + fn get_next(&self) -> Option> { + let main_idx = self.index.load(Ordering::SeqCst); + let provider_count = self.providers.len(); + + for index in main_idx..provider_count { + let priority_group = &self.providers[index]; + let allocation = { + let config = Self::get_next_provider_from_group(priority_group, false); + if config.is_none() { + Self::get_next_provider_from_group(priority_group, true) + } else { + config + } + }; + match allocation { + None => {} + Some(config) => { + if priority_group.is_exhausted() { + self.index.store((index + 1) % provider_count, Ordering::SeqCst); + } + return Some(config); + } + } + } + + None } @@ -497,12 +569,40 @@ impl ActiveProviderManager { Some((lineup, _config)) => lineup.acquire() }; + if log_enabled!(log::Level::Debug) { + match allocation { + ProviderAllocation::Exhausted => {} + ProviderAllocation::Available(ref cfg) | + ProviderAllocation::GracePeriod(ref cfg) => { + debug!("Using provider {}", cfg.name); + } + } + } + ProviderConnectionGuard { manager: Arc::new(self.clone_inner()), allocation, } } + // This method is used for redirects to cycle through provider + // + pub async fn get_next_provider(&self, input_name: &str) -> Option> { + let providers = self.providers.read().await; + match Self::get_provider_config(input_name, &providers) { + None => None, + Some((lineup, _config)) => { + let cfg = lineup.get_next(); + if log_enabled!(log::Level::Debug) { + if let Some(ref c) = cfg { + debug!("Using provider {}", c.name); + } + } + cfg + } + } + } + // we need the provider_name to exactly release this provider pub async fn release_connection(&self, provider_name: &str) { let providers = self.providers.read().await; diff --git a/src/processing/parser/hls.rs b/src/processing/parser/hls.rs index 8721489d6..74be4da8b 100644 --- a/src/processing/parser/hls.rs +++ b/src/processing/parser/hls.rs @@ -12,7 +12,6 @@ pub struct RewriteHlsProps<'a> { pub input_id: u16, } - fn rewrite_hls_url(input: &str, replacement: &str) -> String { if replacement.starts_with('/') { let parts = input.splitn(4, '/').collect::>(); diff --git a/src/processing/processor/playlist.rs b/src/processing/processor/playlist.rs index f9e61c07d..dd3fc5b3d 100644 --- a/src/processing/processor/playlist.rs +++ b/src/processing/processor/playlist.rs @@ -676,7 +676,7 @@ async fn process_playlist_for_target(client: Arc, apply_affixes(&mut processed_fetched_playlists); - let (new_epg, new_playlist) = process_epg(target, &mut processed_fetched_playlists); + let (new_epg, new_playlist) = process_epg(&mut processed_fetched_playlists); if new_playlist.is_empty() { info!("Playlist is empty: {}", &target.name); @@ -691,7 +691,7 @@ async fn process_playlist_for_target(client: Arc, } } -fn process_epg(target: &ConfigTarget, processed_fetched_playlists: &mut Vec) -> (Vec, Vec) { +fn process_epg(processed_fetched_playlists: &mut Vec) -> (Vec, Vec) { let mut new_playlist = vec![]; let mut new_epg = vec![]; @@ -705,7 +705,7 @@ fn process_epg(target: &ConfigTarget, processed_fetched_playlists: &mut Vec