diff --git a/CHANGELOG.md b/CHANGELOG.md index e84d9e7a3..b863e6503 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,13 +1,21 @@ # Changelog # 3.1.0 (2025-05-xx) -- !BREAKING_CHANGE! mapper refactored, mapping can be written as script with a custom DSL. +- !BREAKING_CHANGE! mapper refactored, mapping can be written as a script with a custom DSL. +- !BREAKING_CHANGE! `tags` definition removed from new mapper. +- !BREAKING_CHANGE! removed suffix and prefix from input config. Use mapper with an input filter instead. +- !BREAKING_CHANGE! custom_stream_response is now `custom_stream_response_path`. The filename identifies the file inside the path + - user_account_expired.ts + - provider_connection_exhausted.ts + - user_connection_exhausted.ts + - channel_unavailable.ts + - !BREAKING_CHANGE! epg refactored - url config is now renamed to sources - Added `priority`, priority is `optional` - `auto_epg` is now removed, use `url: auto` instead. - Added `logo_override` to overwrite logo from epg. -**Npte:** The `priority` value determines the importance or order of processing. Lower numbers mean higher priority. That is: +**Note:** The `priority` value determines the importance or order of processing. Lower numbers mean higher priority. That is: A `priority` of `0` is higher than `1`. **Negative numbers** are allowed and represent even higher priority ```yaml diff --git a/README.md b/README.md index 43fd5a434..0abcf3559 100644 --- a/README.md +++ b/README.md @@ -517,8 +517,6 @@ Each input has the following attributes: - `method` can be `GET` or `POST` - `username` only mandatory for type `xtream` - `pasword`only mandatory for type `xtream` -- `prefix` is optional, it is applied to the given field with the given value -- `suffix` is optional, it is applied to the given field with the given value - `options` is optional, + `xtream_skip_live` true or false, live section can be skipped. + `xtream_skip_vod` true or false, vod section can be skipped. @@ -530,11 +528,6 @@ Each input has the following attributes: `persist` should be different for `m3u` and `xtream` types. For `m3u` use full filename like `./playlist_{}.m3u`. For `xtream` use a prefix like `./playlist_` -`prefix` and `suffix` are appended after all processing is done, but before sort. -They have 2 fields: -- `field` can be `name` , `group`, `title`, `caption` -- `value` a static text - Example `epg` config ```yaml @@ -1170,40 +1163,14 @@ mappings: value: '(?i)(?PHD|LQ|4K|UHD)?' - name: source value: 'Url ~ "https?:\/\/(.*?)\/(?P.*)$"' - tags: - - name: quality - captures: - - quality - concat: '|' - prefix: ' [ ' - suffix: ' ]' mapping: - id: France match_as_ascii: true mapper: - filter: 'Name ~ "^TF.*"' - pattern: '!source!' - attributes: - url: http://my.iptv.proxy.com/ - - pattern: 'Name ~ "^TF1$"' - attributes: - name: TF1 - id: TF1.fr - chno: '1' - logo: https://upload.wikimedia.org/wikipedia/commons/thumb/3/3c/TF1_logo_2013.svg/320px-TF1_logo_2013.svg.png - suffix: - title: '' - group: '|FR|TNT' - assignments: - title: name - - pattern: 'Name ~ "^TF1!delimiter!!quality!*Series[_ ]*Films$"' - attributes: - name: TF1 Series Films - id: TF1SeriesFilms.fr - chno: '20' - logo: https://upload.wikimedia.org/wikipedia/commons/thumb/3/3c/TF1_logo_2013.svg/320px-TF1_logo_2013.svg.png, - suffix: - group: '|FR|TNT' + script: | + query_match = @Url ~ "https?:\/\/(.*?)\/(?P.*)$" + @Url = concat("http://my.iptv.proxy.com/", query_match.query) ``` ## 3. Api-Proxy Config diff --git a/bin/build_resources.sh b/bin/build_resources.sh index 90f87f684..bd2f5bfda 100755 --- a/bin/build_resources.sh +++ b/bin/build_resources.sh @@ -21,7 +21,7 @@ while getopts "fh" opt; do esac done -declare -a resources=("channel_unavailable" "user_connections_exhausted" "provider_connections_exhausted") +declare -a resources=("channel_unavailable" "user_connections_exhausted" "provider_connections_exhausted" "user_account_expired") for resource in "${resources[@]}"; do if [ "$flag_force" = false ]; then diff --git a/docker/Dockerfile b/docker/Dockerfile index 887c6da1f..149b9bb62 100644 --- a/docker/Dockerfile +++ b/docker/Dockerfile @@ -61,8 +61,11 @@ RUN ffmpeg -loop 1 -i ./resources/channel_unavailable.jpg -t 10 -r 1 -an \ ffmpeg -loop 1 -i ./resources/provider_connections_exhausted.jpg -t 10 -r 1 -an \ -vf "scale=1920:1080" \ -c:v libx264 -preset veryfast -crf 23 -pix_fmt yuv420p \ - ./resources/provider_connections_exhausted.ts - + ./resources/provider_connections_exhausted.ts && \ + ffmpeg -loop 1 -i ./resources/user_account_expired.jpg -t 10 -r 1 -an \ + -vf "scale=1920:1080" \ + -c:v libx264 -preset veryfast -crf 23 -pix_fmt yuv420p \ + ./resources/user_account_expired.ts # Stage to prepare /etc/localtime for scratch FROM alpine:latest AS tz-prep diff --git a/frontend/src/component/playlist-tree/playlist-tree.tsx b/frontend/src/component/playlist-tree/playlist-tree.tsx index 2bb44bc6b..f9c1eb6bd 100644 --- a/frontend/src/component/playlist-tree/playlist-tree.tsx +++ b/frontend/src/component/playlist-tree/playlist-tree.tsx @@ -40,7 +40,7 @@ export default function PlaylistTree(props: PlaylistTreeProps) { const getPlaylistItemById = useCallback((itemId: string): PlaylistItem => { const id = parseInt(itemId); if (data && !isNaN(id)) { - const groups = [data.live, data.vod, data.series].flat(); + const groups = [data?.live, data?.vod, data?.series].filter(Boolean).flat(); for (let i = 0, len = groups.length; i < len; i++) { const group = groups[i]; for (let j = 0, clen = group.channels?.length ?? 0; j < clen; j++) { diff --git a/resources/channel_unavailable.jpg b/resources/channel_unavailable.jpg index e95ab31d5..7f83bcdfc 100644 Binary files a/resources/channel_unavailable.jpg and b/resources/channel_unavailable.jpg differ diff --git a/resources/provider_connections_exhausted.jpg b/resources/provider_connections_exhausted.jpg index fde4b9228..7a1866346 100644 Binary files a/resources/provider_connections_exhausted.jpg and b/resources/provider_connections_exhausted.jpg differ diff --git a/resources/template.xcf b/resources/template.xcf new file mode 100644 index 000000000..09d00d1ca Binary files /dev/null and b/resources/template.xcf differ diff --git a/resources/user_account_expired.jpg b/resources/user_account_expired.jpg new file mode 100644 index 000000000..f3e247806 Binary files /dev/null and b/resources/user_account_expired.jpg differ diff --git a/resources/user_connections_exhausted.jpg b/resources/user_connections_exhausted.jpg index 3cd076373..5c4a79202 100644 Binary files a/resources/user_connections_exhausted.jpg and b/resources/user_connections_exhausted.jpg differ diff --git a/src/api/api_utils.rs b/src/api/api_utils.rs index 03ebb5a6f..6c8f306c3 100644 --- a/src/api/api_utils.rs +++ b/src/api/api_utils.rs @@ -329,10 +329,10 @@ async fn resolve_streaming_strategy(app_state: &AppState, stream_url: &str, inpu } -fn get_grace_period_millis(connection_permission: &UserConnectionPermission, stream_response_params: &ProviderStreamState, config_grace_period_millis: u64) -> u64 { +fn get_grace_period_millis(connection_permission: UserConnectionPermission, stream_response_params: &ProviderStreamState, config_grace_period_millis: u64) -> u64 { if config_grace_period_millis > 0 && (matches!(stream_response_params, ProviderStreamState::GracePeriod(_, _)) // provider grace period - || connection_permission == &UserConnectionPermission::GracePeriod // user grace period + || connection_permission == UserConnectionPermission::GracePeriod // user grace period ) { config_grace_period_millis } else { 0 } } @@ -350,7 +350,7 @@ async fn create_stream_response_details(app_state: &AppState, resolve_streaming_strategy(app_state, stream_url, input, force_provider).await; let config_grace_period_millis = app_state.config.reverse_proxy.as_ref() .and_then(|r| r.stream.as_ref()).map_or_else(default_grace_period_millis, |s| s.grace_period_millis); - let grace_period_millis = get_grace_period_millis(&connection_permission, &streaming_strategy.provider_stream_state, config_grace_period_millis); + let grace_period_millis = get_grace_period_millis(connection_permission, &streaming_strategy.provider_stream_state, config_grace_period_millis); match streaming_strategy.provider_stream_state { // custom stream means we display our own stream like connection exhausted, channel unavailable... ProviderStreamState::Custom(provider_stream) => { @@ -486,7 +486,7 @@ where 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(params.input, ¶ms.req_context, ¶ms.get_query_path(provider_id, provider_url), provider_url) { + let stream_url = match get_xtream_player_api_stream_url(params.input, params.req_context, ¶ms.get_query_path(provider_id, provider_url), 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()); @@ -567,7 +567,7 @@ pub async fn force_provider_stream_response(app_state: &AppState, let connection_permission = UserConnectionPermission::Allowed; let mut stream_details = - create_stream_response_details(app_state, &stream_options, &user_session.stream_url, req_headers, input, item_type, share_stream, connection_permission.clone(), Some(&user_session.provider)).await; + create_stream_response_details(app_state, &stream_options, &user_session.stream_url, req_headers, input, item_type, share_stream, connection_permission, Some(&user_session.provider)).await; if stream_details.has_stream() { let provider_response = stream_details.stream_info.as_ref().map(|(h, sc)| (h.clone(), *sc)); @@ -608,25 +608,25 @@ pub async fn stream_response(app_state: &AppState, if log_enabled!(log::Level::Trace) { trace!("Try to open stream {}", sanitize_sensitive_info(stream_url)); } if connection_permission == UserConnectionPermission::Exhausted { - return create_custom_video_stream_response(&app_state.config, &CustomVideoStreamType::UserConnectionsExhausted).into_response(); + return create_custom_video_stream_response(&app_state.config, CustomVideoStreamType::UserConnectionsExhausted).into_response(); } let share_stream = is_stream_share_enabled(item_type, target); if share_stream { - if let Some(value) = shared_stream_response(app_state, stream_url, user, connection_permission.clone()).await { + if let Some(value) = shared_stream_response(app_state, stream_url, user, connection_permission).await { return value.into_response(); } } let stream_options = get_stream_options(app_state); let mut stream_details = - create_stream_response_details(app_state, &stream_options, stream_url, req_headers, input, item_type, share_stream, connection_permission.clone(), None).await; + create_stream_response_details(app_state, &stream_options, stream_url, req_headers, input, item_type, share_stream, connection_permission, None).await; if stream_details.has_stream() { // let content_length = get_stream_content_length(provider_response.as_ref()); let provider_response = stream_details.stream_info.as_ref().map(|(h, sc)| (h.clone(), *sc)); let provider_name = stream_details.provider_connection_guard.as_ref().and_then(ProviderConnectionGuard::get_provider_name); - let stream = ActiveClientStream::new(stream_details, app_state, user, connection_permission.clone()).await; + let stream = ActiveClientStream::new(stream_details, app_state, user, connection_permission).await; let stream_resp = if share_stream { debug_if_enabled!("Streaming shared stream request from {}", sanitize_sensitive_info(stream_url)); // Shared Stream response diff --git a/src/api/endpoints/hls_api.rs b/src/api/endpoints/hls_api.rs index d69d3123c..1a849e291 100644 --- a/src/api/endpoints/hls_api.rs +++ b/src/api/endpoints/hls_api.rs @@ -84,7 +84,7 @@ pub(in crate::api) async fn handle_hls_stream_request(app_state: &Arc, } Err(err) => { error!("Failed to download m3u8 {}", sanitize_sensitive_info(err.to_string().as_str())); - create_custom_video_stream_response(&app_state.config, &CustomVideoStreamType::ChannelUnavailable).into_response() + create_custom_video_stream_response(&app_state.config, CustomVideoStreamType::ChannelUnavailable).into_response() } } } @@ -98,7 +98,7 @@ async fn hls_api_stream( app_state.config.get_target_for_user(¶ms.username, ¶ms.password), false, format!("Could not find any user {}", params.username)); if user.permission_denied(&app_state) { - return axum::http::StatusCode::FORBIDDEN.into_response(); + return create_custom_video_stream_response(&app_state.config, CustomVideoStreamType::UserAccountExpired).into_response(); } let target_name = &target.name; @@ -112,11 +112,11 @@ async fn hls_api_stream( if let Some(session) = &mut user_session { if session.permission == UserConnectionPermission::Exhausted { - return create_custom_video_stream_response(&app_state.config, &CustomVideoStreamType::UserConnectionsExhausted).into_response(); + return create_custom_video_stream_response(&app_state.config, CustomVideoStreamType::UserConnectionsExhausted).into_response(); } if app_state.active_provider.is_over_limit(&session.provider).await { - return create_custom_video_stream_response(&app_state.config, &CustomVideoStreamType::ProviderConnectionsExhausted).into_response(); + return create_custom_video_stream_response(&app_state.config, CustomVideoStreamType::ProviderConnectionsExhausted).into_response(); } // let hls_url = if let Some((session_token_opt, hls_url)) = get_hls_session_token_and_url_from_token(&app_state.config.t_encrypt_secret, ¶ms.token) { @@ -147,7 +147,7 @@ async fn hls_api_stream( let connection_permission = user.connection_permission(&app_state).await; if connection_permission == UserConnectionPermission::Exhausted { - return create_custom_video_stream_response(&app_state.config, &CustomVideoStreamType::UserConnectionsExhausted).into_response(); + return create_custom_video_stream_response(&app_state.config, CustomVideoStreamType::UserConnectionsExhausted).into_response(); } if is_hls_url(&session.stream_url) { diff --git a/src/api/endpoints/m3u_api.rs b/src/api/endpoints/m3u_api.rs index becb6fb61..4e12c6d95 100644 --- a/src/api/endpoints/m3u_api.rs +++ b/src/api/endpoints/m3u_api.rs @@ -69,7 +69,7 @@ async fn m3u_api_stream( ) -> impl axum::response::IntoResponse + Send { let (user, target) = try_option_bad_request!(get_user_target_by_credentials(&username, &password, &api_req, &app_state), false, format!("Could not find any user {username}")); if user.permission_denied(&app_state) { - return StatusCode::FORBIDDEN.into_response(); + return create_custom_video_stream_response(&app_state.config, CustomVideoStreamType::UserAccountExpired).into_response(); } let target_name = &target.name; @@ -92,11 +92,11 @@ async fn m3u_api_stream( if let Some(session) = &user_session { if session.permission == UserConnectionPermission::Exhausted { - return create_custom_video_stream_response(&app_state.config, &CustomVideoStreamType::UserConnectionsExhausted).into_response(); + return create_custom_video_stream_response(&app_state.config, CustomVideoStreamType::UserConnectionsExhausted).into_response(); } if app_state.active_provider.is_over_limit(&session.provider).await { - return create_custom_video_stream_response(&app_state.config, &CustomVideoStreamType::ProviderConnectionsExhausted).into_response(); + return create_custom_video_stream_response(&app_state.config, CustomVideoStreamType::ProviderConnectionsExhausted).into_response(); } if session.virtual_id == virtual_id && is_seek_request(cluster, &req_headers).await { // partial request means we are in reverse proxy mode, seek happened @@ -106,7 +106,7 @@ async fn m3u_api_stream( let connection_permission = user.connection_permission(&app_state).await; if connection_permission == UserConnectionPermission::Exhausted { - return create_custom_video_stream_response(&app_state.config, &CustomVideoStreamType::UserConnectionsExhausted).into_response(); + return create_custom_video_stream_response(&app_state.config, CustomVideoStreamType::UserConnectionsExhausted).into_response(); } let context = XtreamApiStreamContext::try_from(cluster).unwrap_or(XtreamApiStreamContext::Live); diff --git a/src/api/endpoints/v1_api.rs b/src/api/endpoints/v1_api.rs index 180cb8338..26748b3d6 100644 --- a/src/api/endpoints/v1_api.rs +++ b/src/api/endpoints/v1_api.rs @@ -261,7 +261,7 @@ async fn config( output: t.output.clone(), rename: t.rename.clone(), mapping: t.mapping.clone(), - processing_order: t.processing_order.clone(), + processing_order: t.processing_order, watch: t.watch.clone(), }; diff --git a/src/api/endpoints/xtream_api.rs b/src/api/endpoints/xtream_api.rs index 000690bbe..302951371 100644 --- a/src/api/endpoints/xtream_api.rs +++ b/src/api/endpoints/xtream_api.rs @@ -42,7 +42,7 @@ use std::path::Path; use std::str::FromStr; use std::sync::Arc; -#[derive(Serialize, Deserialize, Debug, Clone, Eq, PartialEq)] +#[derive(Serialize, Deserialize, Debug, Copy, Clone, Eq, PartialEq)] pub enum XtreamApiStreamContext { LiveAlt, Live, @@ -140,7 +140,7 @@ pub fn serve_query(file_path: &Path, filter: &HashMap<&str, HashSet>) -> } pub(in crate::api) fn get_xtream_player_api_stream_url( - input: &ConfigInput, context: &XtreamApiStreamContext, action_path: &str, fallback_url: &str, + input: &ConfigInput, context: XtreamApiStreamContext, action_path: &str, fallback_url: &str, ) -> Option { if let Some(input_user_info) = input.get_user_info() { let ctx = match context { @@ -182,7 +182,7 @@ async fn xtream_player_api_stream( ) -> impl IntoResponse + Send { 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)); if user.permission_denied(app_state) { - return StatusCode::FORBIDDEN.into_response(); + return create_custom_video_stream_response(&app_state.config, CustomVideoStreamType::UserAccountExpired).into_response(); } let target_name = &target.name; @@ -204,11 +204,11 @@ async fn xtream_player_api_stream( if let Some(session) = &user_session { if session.permission == UserConnectionPermission::Exhausted { - return create_custom_video_stream_response(&app_state.config, &CustomVideoStreamType::UserConnectionsExhausted).into_response(); + return create_custom_video_stream_response(&app_state.config, CustomVideoStreamType::UserConnectionsExhausted).into_response(); } if app_state.active_provider.is_over_limit(&session.provider).await { - return create_custom_video_stream_response(&app_state.config, &CustomVideoStreamType::ProviderConnectionsExhausted).into_response(); + return create_custom_video_stream_response(&app_state.config, CustomVideoStreamType::ProviderConnectionsExhausted).into_response(); } if session.virtual_id == virtual_id && is_seek_request(cluster, req_headers).await { @@ -219,10 +219,10 @@ async fn xtream_player_api_stream( let connection_permission = user.connection_permission(app_state).await; if connection_permission == UserConnectionPermission::Exhausted { - return create_custom_video_stream_response(&app_state.config, &CustomVideoStreamType::UserConnectionsExhausted).into_response(); + return create_custom_video_stream_response(&app_state.config, CustomVideoStreamType::UserConnectionsExhausted).into_response(); } - let context = stream_req.context.clone(); + let context = stream_req.context; let redirect_params = RedirectParams { item: &pli, @@ -249,7 +249,7 @@ async fn xtream_player_api_stream( 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), + let stream_url = try_option_bad_request!(get_xtream_player_api_stream_url(input, 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 || pli.item_type == PlaylistItemType::LiveDash || extension == HLS_EXT; @@ -314,7 +314,7 @@ async fn xtream_player_api_stream_with_token( }; 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.as_str()), true, format!("Cant find stream url for target {target_name}, context {}, stream_id {virtual_id}", stream_req.context)); diff --git a/src/api/model/active_provider_manager.rs b/src/api/model/active_provider_manager.rs index 3ac47eaf1..eb1b50c44 100644 --- a/src/api/model/active_provider_manager.rs +++ b/src/api/model/active_provider_manager.rs @@ -65,7 +65,7 @@ impl Drop for ProviderConnectionGuard { } } -#[derive(Debug)] +#[derive(Debug, Clone)] pub enum ProviderAllocation { Exhausted, Available(Arc), @@ -621,8 +621,6 @@ mod tests { username: None, password: None, persist: None, - prefix: None, - suffix: None, enabled: true, input_type: InputType::Xtream, // You can use a default value here max_connections, diff --git a/src/api/model/active_user_manager.rs b/src/api/model/active_user_manager.rs index c54787111..54040240e 100644 --- a/src/api/model/active_user_manager.rs +++ b/src/api/model/active_user_manager.rs @@ -250,7 +250,7 @@ impl ActiveUserManager { } if let Some(index) = found_session_index { - let session_permission = connection_data.sessions[index].permission.clone(); + let session_permission = connection_data.sessions[index].permission; if session_permission == UserConnectionPermission::GracePeriod { let new_permission = self.check_connection_permission(username, connection_data); connection_data.sessions[index].permission = new_permission; diff --git a/src/api/model/provider_config.rs b/src/api/model/provider_config.rs index c86bcc92d..8ad008b72 100644 --- a/src/api/model/provider_config.rs +++ b/src/api/model/provider_config.rs @@ -6,7 +6,7 @@ use std::ops::Deref; use std::sync::Arc; use tokio::sync::RwLock; -#[derive(Debug)] +#[derive(Debug, Clone, Copy)] pub enum ProviderConfigAllocation { Exhausted, Available, diff --git a/src/api/model/streams/active_client_stream.rs b/src/api/model/streams/active_client_stream.rs index 13459eac4..f68d3028f 100644 --- a/src/api/model/streams/active_client_stream.rs +++ b/src/api/model/streams/active_client_stream.rs @@ -1,12 +1,12 @@ -use crate::api::model::active_user_manager::UserConnectionGuard; use crate::api::api_utils::StreamDetails; use crate::api::model::active_provider_manager::{ActiveProviderManager, ProviderConnectionGuard}; use crate::api::model::active_user_manager::ActiveUserManager; +use crate::api::model::active_user_manager::UserConnectionGuard; use crate::api::model::app_state::AppState; use crate::api::model::stream::BoxedProviderStream; use crate::api::model::stream_error::StreamError; use crate::api::model::streams::chunked_buffer::ChunkedBuffer; -use crate::model::{ProxyUserCredentials, UserConnectionPermission}; +use crate::model::{CustomStreamResponse, ProxyUserCredentials, UserConnectionPermission}; use bytes::Bytes; use futures::Stream; use log::{error, info}; @@ -24,7 +24,8 @@ pub(in crate::api) struct ActiveClientStream { #[allow(unused)] user_connection_guard: Option, send_custom_stream_flag: Option>, - custom_video: (Option, Option), + custom_video: Option, + active_custom_video: Option, #[allow(dead_code)] provider_connection_guard: Option, } @@ -48,10 +49,8 @@ impl ActiveClientStream { inner: stream_details.stream.take().unwrap(), user_connection_guard, send_custom_stream_flag: grace_stop_flag, - custom_video: ( - config.t_provider_connections_exhausted_video.as_ref().map(|a| ChunkedBuffer::new(Arc::clone(a))), - config.t_user_connections_exhausted_video.as_ref().map(|a| ChunkedBuffer::new(Arc::clone(a))), - ), + custom_video: config.t_custom_stream_response.clone(), + active_custom_video: None, provider_connection_guard: stream_details.provider_connection_guard, } } @@ -108,7 +107,6 @@ impl ActiveClientStream { } None } - } impl Stream for ActiveClientStream { type Item = Result; @@ -123,10 +121,27 @@ impl Stream for ActiveClientStream { if self.provider_connection_guard.is_some() { drop(self.provider_connection_guard.take()); } - let custom_video = if send_custom_stream == PROVIDER_EXHAUSTED_FLAG { - self.custom_video.0.as_mut() + + if self.active_custom_video.is_none() { + let video = match self.custom_video.as_ref() { + None => None, + Some(custom_videos) => { + if send_custom_stream == PROVIDER_EXHAUSTED_FLAG { + custom_videos.provider_connections_exhausted.as_ref() + } else { + custom_videos.user_connections_exhausted.as_ref() + } + } + }; + if let Some(custom_video) = video.as_ref() { + self.active_custom_video = Some(ChunkedBuffer::new(Arc::clone(custom_video))); + } + } + + let custom_video: Option<&mut ChunkedBuffer> = if self.active_custom_video.is_some() { + self.active_custom_video.as_mut() } else { - self.custom_video.1.as_mut() + None }; return match custom_video { None => { diff --git a/src/api/model/streams/provider_stream.rs b/src/api/model/streams/provider_stream.rs index 740c67918..dd9ed616e 100644 --- a/src/api/model/streams/provider_stream.rs +++ b/src/api/model/streams/provider_stream.rs @@ -8,10 +8,12 @@ use std::sync::Arc; use axum::response::IntoResponse; use crate::api::model::stream::ProviderStreamResponse; +#[derive(Debug, Copy, Clone)] pub enum CustomVideoStreamType { ChannelUnavailable, UserConnectionsExhausted, ProviderConnectionsExhausted, + UserAccountExpired } fn create_video_stream(video: Option<&Arc>>, headers: &[(String, String)], log_message: &str) -> ProviderStreamResponse { @@ -28,22 +30,31 @@ fn create_video_stream(video: Option<&Arc>>, headers: &[(String, String) } pub fn create_channel_unavailable_stream(cfg: &Config, headers: &[(String, String)], status: StatusCode) -> ProviderStreamResponse { - create_video_stream(cfg.t_channel_unavailable_video.as_ref(), headers, &format!("Streaming response channel unavailable for status {status}")) + let video = cfg.t_custom_stream_response.as_ref().and_then(|c| c.channel_unavailable.as_ref()); + create_video_stream(video, headers, &format!("Streaming response channel unavailable for status {status}")) } pub fn create_user_connections_exhausted_stream(cfg: &Config, headers: &[(String, String)]) -> ProviderStreamResponse { - create_video_stream(cfg.t_user_connections_exhausted_video.as_ref(), headers, "Streaming response user connections exhausted") + let video = cfg.t_custom_stream_response.as_ref().and_then(|c| c.user_connections_exhausted.as_ref()); + create_video_stream(video, headers, "Streaming response user connections exhausted") } pub fn create_provider_connections_exhausted_stream(cfg: &Config, headers: &[(String, String)]) -> ProviderStreamResponse { - create_video_stream(cfg.t_provider_connections_exhausted_video.as_ref(), headers, "Streaming response provider connections exhausted") + let video = cfg.t_custom_stream_response.as_ref().and_then(|c| c.provider_connections_exhausted.as_ref()); + create_video_stream(video, headers, "Streaming response provider connections exhausted") } -pub fn create_custom_video_stream_response(config: &Config, video_response: &CustomVideoStreamType) -> impl axum::response::IntoResponse + Send { +pub fn create_user_account_expired_stream(cfg: &Config, headers: &[(String, String)]) -> ProviderStreamResponse { + let video = cfg.t_custom_stream_response.as_ref().and_then(|c| c.user_account_expired.as_ref()); + create_video_stream(video, headers, "Streaming response user account expired") +} + +pub fn create_custom_video_stream_response(config: &Config, video_response: CustomVideoStreamType) -> impl axum::response::IntoResponse + Send { if let (Some(stream), Some((headers, status_code))) = match video_response { CustomVideoStreamType::ChannelUnavailable => create_channel_unavailable_stream(config, &[], StatusCode::BAD_REQUEST), CustomVideoStreamType::UserConnectionsExhausted => create_user_connections_exhausted_stream(config, &[]), CustomVideoStreamType::ProviderConnectionsExhausted => create_provider_connections_exhausted_stream(config, &[]), + CustomVideoStreamType::UserAccountExpired => create_user_account_expired_stream(config, &[]), } { let mut builder = axum::response::Response::builder() .status(status_code); diff --git a/src/foundation/filter.rs b/src/foundation/filter.rs index 908b5beef..ab0b14ba2 100644 --- a/src/foundation/filter.rs +++ b/src/foundation/filter.rs @@ -16,7 +16,7 @@ use crate::tuliprox_error::{create_tuliprox_error_result, info_err}; use crate::tuliprox_error::{TuliproxError, TuliproxErrorKind}; use crate::utils::CONSTANTS; -pub fn get_field_value(pli: &PlaylistItem, field: &ItemField) -> String { +pub fn get_field_value(pli: &PlaylistItem, field: ItemField) -> String { let header = &pli.header; let value = match field { ItemField::Group => header.group.to_string(), @@ -30,7 +30,7 @@ pub fn get_field_value(pli: &PlaylistItem, field: &ItemField) -> String { value.to_string() } -pub fn set_field_value(pli: &mut PlaylistItem, field: &ItemField, value: String) -> bool { +pub fn set_field_value(pli: &mut PlaylistItem, field: ItemField, value: String) -> bool { let header = &mut pli.header; match field { ItemField::Group => header.group = value, @@ -134,12 +134,12 @@ main = _{ SOI ~ stmt ~ EOI } "#] struct FilterParser; -#[derive(Debug, Clone)] +#[derive(Debug, Copy, Clone)] pub enum UnaryOperator { Not, } -#[derive(Debug, Clone)] +#[derive(Debug, Copy, Clone)] pub enum BinaryOperator { And, Or, @@ -478,7 +478,7 @@ pub fn get_filter( Some(binop) => { result = Some(Filter::BinaryExpression( Box::new(result.unwrap()), - binop.clone(), + *binop, Box::new(expr), )); op = None; diff --git a/src/foundation/mapper.rs b/src/foundation/mapper.rs index 0aa39435d..48ed63209 100644 --- a/src/foundation/mapper.rs +++ b/src/foundation/mapper.rs @@ -1005,7 +1005,7 @@ mod tests { ("K", "SHD"), ("C", "LHD"), ("L", "FHD"), ("R", "UHD"), ("T", "SD"), ("A", "FHD"), ].into_iter().map(|(name, quality)| PlaylistItem { header: PlaylistItemHeader { title: format!("Chanel {name} [{quality}]"), ..Default::default() } }).collect::>(); - for pli in channels.iter_mut() { + for pli in &mut channels { let mut accessor = ValueAccessor { pli, }; diff --git a/src/messaging.rs b/src/messaging.rs index d60da306c..5b4fb8b5a 100644 --- a/src/messaging.rs +++ b/src/messaging.rs @@ -3,7 +3,7 @@ use crate::model::{MessagingConfig}; use log::{debug, error}; use reqwest::{header}; -#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, PartialEq, Eq)] +#[derive(Debug, Copy, Clone, serde::Serialize, serde::Deserialize, PartialEq, Eq)] pub enum MsgKind { #[serde(rename = "info")] Info, @@ -15,8 +15,8 @@ pub enum MsgKind { Watch, } -fn is_enabled(kind: &MsgKind, cfg: &MessagingConfig) -> bool { - cfg.notify_on.contains(kind) +fn is_enabled(kind: MsgKind, cfg: &MessagingConfig) -> bool { + cfg.notify_on.contains(&kind) } fn send_http_post_request(client: &Arc, msg: &str, messaging: &MessagingConfig) { @@ -84,7 +84,7 @@ fn send_pushover_message(client: &Arc, msg: &str, messaging: &M pub fn send_message(client: &Arc, kind: &MsgKind, cfg: Option<&MessagingConfig>, msg: &str) { if let Some(messaging) = cfg { - if is_enabled(kind, messaging) { + if is_enabled(*kind, messaging) { send_telegram_message(msg, messaging); send_http_post_request(client, msg, messaging); send_pushover_message(client, msg, messaging); diff --git a/src/model/api_proxy.rs b/src/model/api_proxy.rs index 2c5e39bb9..d5c6d151f 100644 --- a/src/model/api_proxy.rs +++ b/src/model/api_proxy.rs @@ -15,7 +15,7 @@ use serde::{Deserialize, Deserializer, Serialize, Serializer}; use crate::model::PlaylistItemType; use crate::utils; -#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, PartialEq, Eq)] +#[derive(Debug, Copy, Clone, serde::Serialize, serde::Deserialize, PartialEq, Eq)] pub enum UserConnectionPermission { Exhausted, Allowed, @@ -130,7 +130,7 @@ impl Serialize for ProxyType { } -#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, Sequence, PartialEq, Eq)] +#[derive(Debug, Copy, Clone, serde::Serialize, serde::Deserialize, Sequence, PartialEq, Eq)] pub enum ProxyUserStatus { Active, // The account is in good standing and can stream content Expired, // The account can no longer access content unless it is renewed. diff --git a/src/model/config.rs b/src/model/config.rs index 74ffc8c52..f29e3ba5a 100644 --- a/src/model/config.rs +++ b/src/model/config.rs @@ -3,7 +3,7 @@ use arc_swap::ArcSwapOption; use enum_iterator::Sequence; use std::collections::{HashMap, HashSet}; use std::fmt::Display; -use std::path::PathBuf; +use std::path::{Path, PathBuf}; use std::str::FromStr; use std::sync::Arc; @@ -39,7 +39,7 @@ macro_rules! valid_property { pub use valid_property; use crate::utils; -#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, Sequence, Eq, PartialEq)] +#[derive(Debug, Copy, Clone, serde::Serialize, serde::Deserialize, Sequence, Eq, PartialEq)] pub enum ItemField { #[serde(rename = "group")] Group, @@ -118,7 +118,7 @@ impl FromStr for ItemField { } -#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] +#[derive(Debug, Copy, Clone, serde::Serialize, serde::Deserialize)] pub enum FilterMode { #[serde(rename = "discard")] Discard, @@ -361,13 +361,15 @@ impl ReverseProxyConfig { #[derive(Debug, Clone, serde::Serialize, serde::Deserialize, Default)] #[serde(deny_unknown_fields)] -pub struct CustomStreamResponseConfig { - #[serde(default, skip_serializing_if = "Option::is_none")] - pub channel_unavailable: Option, - #[serde(default, skip_serializing_if = "Option::is_none")] - pub user_connections_exhausted: Option, // user has no more connections - #[serde(default, skip_serializing_if = "Option::is_none")] - pub provider_connections_exhausted: Option, // provider limit reached, has no more connections +pub struct CustomStreamResponse { + #[serde(default, skip)] + pub channel_unavailable: Option>>, + #[serde(default, skip)] + pub user_connections_exhausted: Option>>, // user has no more connections + #[serde(default, skip)] + pub provider_connections_exhausted: Option>>, // provider limit reached, has no more connections + #[serde(default, skip)] + pub user_account_expired: Option>>, } @@ -386,7 +388,7 @@ pub struct Config { #[serde(default, skip_serializing_if = "Option::is_none")] pub mapping_path: Option, #[serde(default, skip_serializing_if = "Option::is_none")] - pub custom_stream_response: Option, + pub custom_stream_response_path: Option, #[serde(default, skip_serializing_if = "Option::is_none")] pub templates: Option>, #[serde(default, skip_serializing_if = "Option::is_none")] @@ -432,11 +434,7 @@ pub struct Config { #[serde(skip)] pub file_locks: Arc, #[serde(skip)] - pub t_channel_unavailable_video: Option>>, - #[serde(skip)] - pub t_user_connections_exhausted_video: Option>>, - #[serde(skip)] - pub t_provider_connections_exhausted_video: Option>>, + pub t_custom_stream_response: Option, #[serde(skip)] pub t_access_token_secret: [u8; 32], #[serde(skip)] @@ -801,22 +799,33 @@ impl Config { } fn prepare_custom_stream_response(&mut self) { - if let Some(custom_stream_response) = self.custom_stream_response.as_ref() { - fn load_and_set_file(path: Option<&String>, working_dir: &str) -> Option>> { - path.as_ref() - .map(|file| utils::make_absolute_path(file, working_dir)) - .and_then(|absolute_path| match utils::read_file_as_bytes(&PathBuf::from(&absolute_path)) { + if let Some(custom_stream_response_path) = self.custom_stream_response_path.as_ref() { + fn load_and_set_file(path: &Path, working_dir: &str) -> Option>> { + let file_path = utils::make_path_absolute(path, working_dir); + if file_path.exists() { + match utils::read_file_as_bytes(&PathBuf::from(&file_path)) { Ok(data) => Some(Arc::new(data)), Err(err) => { - error!("Failed to load file: {absolute_path} {err}"); + error!("Failed to load a resource file: {} {err}", file_path.display()); None } - }) + } + } else { + None + } } - self.t_channel_unavailable_video = load_and_set_file(custom_stream_response.channel_unavailable.as_ref(), &self.working_dir); - self.t_user_connections_exhausted_video = load_and_set_file(custom_stream_response.user_connections_exhausted.as_ref(), &self.working_dir); - self.t_provider_connections_exhausted_video = load_and_set_file(custom_stream_response.provider_connections_exhausted.as_ref(), &self.working_dir); + let path = PathBuf::from(custom_stream_response_path); + let channel_unavailable = load_and_set_file(&path.join("channel_unavailable.ts"), &self.working_dir); + let user_connections_exhausted = load_and_set_file(&path.join("user_connections_exhausted.ts"), &self.working_dir); + let provider_connections_exhausted = load_and_set_file(&path.join("provider_connections_exhausted.ts"), &self.working_dir); + let user_account_expired = load_and_set_file(&path.join("user_account_expired.ts"), &self.working_dir); + self.t_custom_stream_response = Some(CustomStreamResponse { + channel_unavailable, + user_connections_exhausted, + provider_connections_exhausted, + user_account_expired, + }); } } diff --git a/src/model/config_input.rs b/src/model/config_input.rs index 864a67f2f..90daef685 100644 --- a/src/model/config_input.rs +++ b/src/model/config_input.rs @@ -40,17 +40,7 @@ pub struct InputAffix { pub value: String, } -#[derive( - Debug, - Copy, - Clone, - serde::Serialize, - serde::Deserialize, - Sequence, - PartialEq, - Eq, - Default -)] +#[derive(Debug,Copy,Clone,serde::Serialize,serde::Deserialize,Sequence,PartialEq,Eq,Default)] pub enum InputType { #[serde(rename = "m3u")] #[default] @@ -99,17 +89,7 @@ impl FromStr for InputType { } } -#[derive( - Debug, - Copy, - Clone, - serde::Serialize, - serde::Deserialize, - Sequence, - PartialEq, - Eq, - Default -)] +#[derive(Debug, Copy, Clone, serde::Serialize, serde::Deserialize, Sequence, PartialEq, Eq, Default)] pub enum InputFetchMethod { #[default] GET, @@ -253,10 +233,6 @@ pub struct ConfigInput { pub password: Option, #[serde(default, skip_serializing_if = "Option::is_none")] pub persist: Option, - #[serde(default, skip_serializing_if = "Option::is_none")] - pub prefix: Option, - #[serde(default, skip_serializing_if = "Option::is_none")] - pub suffix: Option, #[serde(default = "default_as_true")] pub enabled: bool, #[serde(default, skip_serializing_if = "Option::is_none")] diff --git a/src/model/config_sort.rs b/src/model/config_sort.rs index e7915bc6f..76d4d7c75 100644 --- a/src/model/config_sort.rs +++ b/src/model/config_sort.rs @@ -3,7 +3,7 @@ use crate::foundation::filter::{apply_templates_to_pattern, apply_templates_to_p use crate::tuliprox_error::{TuliproxError, TuliproxErrorKind, create_tuliprox_error, handle_tuliprox_error_result_list}; use crate::model::{ItemField}; -#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] +#[derive(Debug, Copy, Clone, serde::Serialize, serde::Deserialize)] pub enum SortOrder { #[serde(rename = "asc")] Asc, diff --git a/src/model/mapping.rs b/src/model/mapping.rs index 17b5418ad..53f347898 100644 --- a/src/model/mapping.rs +++ b/src/model/mapping.rs @@ -18,7 +18,7 @@ pub const MAPPER_FIELDS: &[&str] = &[ "time_shift", "rec", "url", "epg_channel_id", "epg_id" ]; -#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, Sequence, PartialEq, Eq)] +#[derive(Debug, Copy, Clone, serde::Serialize, serde::Deserialize, Sequence, PartialEq, Eq)] pub enum CounterModifier { #[serde(rename = "assign")] Assign, @@ -200,7 +200,7 @@ impl Mapping { filter: flt, field: def.field.clone(), concat: def.concat.clone(), - modifier: def.modifier.clone(), + modifier: def.modifier, value: Arc::new(AtomicU32::new(def.value)), padding: def.padding, }); diff --git a/src/model/processing_order.rs b/src/model/processing_order.rs index ae66ab98d..79feaa0a4 100644 --- a/src/model/processing_order.rs +++ b/src/model/processing_order.rs @@ -1,7 +1,7 @@ use std::fmt::Display; use enum_iterator::Sequence; -#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, Sequence, PartialEq, Eq, Default)] +#[derive(Debug, Copy, Clone, serde::Serialize, serde::Deserialize, Sequence, PartialEq, Eq, Default)] pub enum ProcessingOrder { #[serde(rename = "frm")] #[default] diff --git a/src/processing/processor/playlist.rs b/src/processing/processor/playlist.rs index 439c0dd7c..f2e372ad3 100644 --- a/src/processing/processor/playlist.rs +++ b/src/processing/processor/playlist.rs @@ -81,13 +81,13 @@ fn exec_rename(pli: &mut PlaylistItem, rename: Option<&Vec>) { if !renames.is_empty() { let result = pli; for r in renames { - let value = get_field_value(result, &r.field); + let value = get_field_value(result, r.field); let cap = r.re.as_ref().unwrap().replace_all(value.as_str(), &r.new_name); if log_enabled!(log::Level::Debug) && *value != cap { debug_if_enabled!("Renamed {}={} to {}", &r.field, value, cap); } let value = cap.into_owned(); - set_field_value(result, &r.field, value); + set_field_value(result, r.field, value); } } } diff --git a/src/processing/processor/sort.rs b/src/processing/processor/sort.rs index 2cae63eec..4b0a28155 100644 --- a/src/processing/processor/sort.rs +++ b/src/processing/processor/sort.rs @@ -6,7 +6,7 @@ use std::cmp::Ordering; fn playlist_comparator( sequence: Option<&Vec>, - order: &SortOrder, + order: SortOrder, value_a: &str, value_b: &str, ) -> Ordering { @@ -98,7 +98,7 @@ fn playlistgroup_comparator(a: &PlaylistGroup, b: &PlaylistGroup, group_sort: &C let value_a = if match_as_ascii { deunicode(&a.title) } else { a.title.to_string() }; let value_b = if match_as_ascii { deunicode(&b.title) } else { b.title.to_string() }; - playlist_comparator(group_sort.t_re_sequence.as_ref(), &group_sort.order, &value_a, &value_b) + playlist_comparator(group_sort.t_re_sequence.as_ref(), group_sort.order, &value_a, &value_b) } fn playlistitem_comparator( @@ -107,12 +107,12 @@ fn playlistitem_comparator( channel_sort: &ConfigSortChannel, match_as_ascii: bool, ) -> Ordering { - let raw_value_a = get_field_value(a, &channel_sort.field); - let raw_value_b = get_field_value(b, &channel_sort.field); + let raw_value_a = get_field_value(a, channel_sort.field); + let raw_value_b = get_field_value(b, channel_sort.field); let value_a = if match_as_ascii { deunicode(&raw_value_a) } else { raw_value_a }; let value_b = if match_as_ascii { deunicode(&raw_value_b) } else { raw_value_b }; - playlist_comparator(channel_sort.t_re_sequence.as_ref(), &channel_sort.order, &value_a, &value_b) + playlist_comparator(channel_sort.t_re_sequence.as_ref(), channel_sort.order, &value_a, &value_b) } pub(in crate::processing::processor) fn sort_playlist(target: &ConfigTarget, new_playlist: &mut [PlaylistGroup]) { diff --git a/src/repository/user_repository.rs b/src/repository/user_repository.rs index cd2caf00f..db18bd930 100644 --- a/src/repository/user_repository.rs +++ b/src/repository/user_repository.rs @@ -42,7 +42,7 @@ impl StoredProxyUserCredentialsDeprecated { created_at: stored.created_at, exp_date: stored.exp_date, max_connections: stored.max_connections.unwrap_or_default(), - status: stored.status.clone(), + status: stored.status, ui_enabled: stored.ui_enabled, comment: None, } @@ -83,7 +83,7 @@ impl StoredProxyUserCredentials { created_at: proxy.created_at, exp_date: proxy.exp_date, max_connections: if proxy.max_connections > 0 { Some(proxy.max_connections) } else { None }, - status: proxy.status.clone(), + status: proxy.status, ui_enabled: proxy.ui_enabled, comment: proxy.comment.clone(), } @@ -100,7 +100,7 @@ impl StoredProxyUserCredentials { created_at: stored.created_at, exp_date: stored.exp_date, max_connections: stored.max_connections.unwrap_or_default(), - status: stored.status.clone(), + status: stored.status, ui_enabled: stored.ui_enabled, comment: stored.comment.clone(), } diff --git a/src/tuliprox_error.rs b/src/tuliprox_error.rs index c1fb8dc05..2aedf8d6f 100644 --- a/src/tuliprox_error.rs +++ b/src/tuliprox_error.rs @@ -90,7 +90,7 @@ macro_rules! handle_tuliprox_error_result { pub use handle_tuliprox_error_result; -#[derive(Debug, PartialEq, Eq)] +#[derive(Debug, Copy, Clone, PartialEq, Eq)] pub enum TuliproxErrorKind { // do not send with messaging Info, diff --git a/src/utils/file/file_utils.rs b/src/utils/file/file_utils.rs index 95ac0eacc..ae5606b01 100644 --- a/src/utils/file/file_utils.rs +++ b/src/utils/file/file_utils.rs @@ -271,22 +271,27 @@ pub fn read_file_as_bytes(path: &Path) -> std::io::Result> { pub fn make_absolute_path(path: &str, working_dir: &str) -> String { let rpb = std::path::PathBuf::from(path); + let pathbuf = make_path_absolute(&rpb, working_dir); + pathbuf.to_str().unwrap_or_default().to_string() +} + +pub fn make_path_absolute(rpb: &Path, working_dir: &str) -> PathBuf { if rpb.is_relative() { - let mut rpb2 = std::path::PathBuf::from(working_dir).join(&rpb); + let mut rpb2 = std::path::PathBuf::from(working_dir).join(rpb); if !rpb2.exists() { - rpb2 = get_exe_path().join(&rpb); + rpb2 = get_exe_path().join(rpb); } if !rpb2.exists() { let cwd = std::env::current_dir(); if let Ok(cwd_path) = cwd { - rpb2 = cwd_path.join(&rpb); + rpb2 = cwd_path.join(rpb); } } if rpb2.exists() { - return String::from(rpb2.clean().to_str().unwrap_or_default()); + return rpb2.clean(); } } - path.to_string() + rpb.to_path_buf() } pub fn resolve_relative_path(relative: &str) -> std::io::Result {