diff --git a/README.md b/README.md index a5a1491c5..2cece4210 100644 --- a/README.md +++ b/README.md @@ -789,7 +789,7 @@ Each input has the following attributes: - `username` only mandatory for type `xtream` - `password` only mandatory for type `xtream` - `panel_api` _optional_ for provider panel api operations -- `exp_date` optional, i a date as "YYYY-MM-DD HH:MM:SS" format like `2028-11-30 12:34:12` or Unix timestamp (seconds since epoch) +- `exp_date` optional, is a date as "YYYY-MM-DD HH:MM:SS" format like `2028-11-30 12:34:12` or Unix timestamp (seconds since epoch) - `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. @@ -977,6 +977,8 @@ If provider connections are exhausted, tuliprox can optionally call a provider p - renew expired accounts first (based on `exp_date`) - otherwise create a new alias account and persist it +**Important!** Panel api accounts are not considering unlimited provider access! + The API is configured generically via predefined query parameters; only `type: m3u` is supported. Use the literal value `auto` to fill sensitive values at runtime: @@ -1001,13 +1003,13 @@ Tuliprox evaluates Panel API responses as JSON with the following logic, dependi `client_new (create alias)` -- Require `status: true`. -- Attempt to extract credentials directly from the JSON response: - - username - - password -- If one or both fields are missing, tuliprox attempts a fallback extraction from a URL contained in the JSON: - - If the JSON contains a url field, tuliprox parses it and tries to extract username/password from it (e.g., query string or embedded credentials depending on the provider’s URL format). -- If credentials cannot be derived from either the direct fields or the url fallback, the operation is treated as failed and no alias is persisted. +- Require `status: true`. +- Attempt to extract credentials directly from the JSON response: + - username + - password +- If one or both fields are missing, tuliprox attempts a fallback extraction from a URL contained in the JSON: + - If the JSON contains a url field, tuliprox parses it and tries to extract username/password from it (e.g., query string or embedded credentials depending on the provider’s URL format). +- If credentials cannot be derived from either the direct fields or the url fallback, the operation is treated as failed and no alias is persisted. `client_renew (renew existing account)` @@ -1107,7 +1109,6 @@ Channels within the `Freetv` group are first sorted by `quality` (as matched by To sort by specific parts of the content, use named capture groups such as `c1`, `c2`, `c3`, etc. The numeric suffix indicates the priority: `c1` is evaluated first, followed by `c2`, and so on. - ### 2.2.2.2 `output` Is a list of output format: diff --git a/backend/src/api/model/provider_config.rs b/backend/src/api/model/provider_config.rs index 3d0d5eac8..667b21bfe 100644 --- a/backend/src/api/model/provider_config.rs +++ b/backend/src/api/model/provider_config.rs @@ -103,6 +103,7 @@ macro_rules! modify_connections { impl ProviderConfig { pub fn new(cfg: &ConfigInput, connection: Arc>, on_connection_change: ProviderConnectionChangeCallback) -> Self { let panel_api_enabled = cfg.panel_api.as_ref().is_some_and(|panel_api| panel_api.enabled); + // Logic change: panel api accounts are not considering unlimited provider access! let effective_max_connections = if panel_api_enabled && cfg.max_connections == 0 { debug_if_enabled!( "panel_api: input '{}' has max_connections=0; defaulting effective max_connections to 1 for pool accounting", diff --git a/backend/src/api/model/streams/active_client_stream.rs b/backend/src/api/model/streams/active_client_stream.rs index dbed4df3d..49a5e06cb 100644 --- a/backend/src/api/model/streams/active_client_stream.rs +++ b/backend/src/api/model/streams/active_client_stream.rs @@ -1,24 +1,24 @@ -use crate::api::model::{AppState, ConnectionManager, CustomVideoStreamType, ProviderHandle, StreamDetails}; use crate::api::model::BoxedProviderStream; use crate::api::model::StreamError; use crate::api::model::TimedClientStream; use crate::api::model::TransportStreamBuffer; -use crate::utils::debug_if_enabled; +use crate::api::model::{AppState, ConnectionManager, CustomVideoStreamType, ProviderHandle, StreamDetails}; +use crate::auth::Fingerprint; use crate::model::ProxyUserCredentials; +use crate::utils::debug_if_enabled; +use axum::http::header::USER_AGENT; +use axum::http::HeaderMap; use bytes::Bytes; +use futures::task::AtomicWaker; use futures::Stream; use futures::StreamExt; use log::{error, info}; use shared::model::{StreamChannel, UserConnectionPermission}; +use shared::utils::sanitize_sensitive_info; use std::pin::Pin; use std::sync::atomic::AtomicU8; -use std::sync::{Arc}; -use std::task::{Poll}; -use axum::http::header::USER_AGENT; -use axum::http::HeaderMap; -use futures::task::AtomicWaker; -use shared::utils::sanitize_sensitive_info; -use crate::auth::Fingerprint; +use std::sync::Arc; +use std::task::Poll; const INNER_STREAM: u8 = 0_u8; const USER_EXHAUSTED_STREAM: u8 = 1_u8; @@ -34,19 +34,19 @@ pub(in crate::api) struct ActiveClientStream { waker: Option>, connection_manager: Arc, fingerprint: Arc, + provider_stopped: bool, } impl ActiveClientStream { - #[allow(clippy::too_many_arguments)] pub(crate) async fn new(mut stream_details: StreamDetails, - app_state: &Arc, - user: &ProxyUserCredentials, - connection_permission: UserConnectionPermission, - fingerprint: &Fingerprint, - stream_channel: StreamChannel, - session_token: Option<&str>, - req_headers: &HeaderMap) -> Self { + app_state: &Arc, + user: &ProxyUserCredentials, + connection_permission: UserConnectionPermission, + fingerprint: &Fingerprint, + stream_channel: StreamChannel, + session_token: Option<&str>, + req_headers: &HeaderMap) -> Self { if connection_permission == UserConnectionPermission::Exhausted { error!("Something is wrong this should not happen"); } @@ -58,14 +58,19 @@ impl ActiveClientStream { let virtual_id = stream_channel.virtual_id; app_state.connection_manager.update_connection(username, user.max_connections, fingerprint, &provider_name, stream_channel, user_agent, session_token).await; - if let Some((_,_,_m_, Some(cvt))) = stream_details.stream_info.as_ref() { + if let Some((_, _, _m_, Some(cvt))) = stream_details.stream_info.as_ref() { app_state.connection_manager.update_stream_detail(&fingerprint.addr, *cvt).await; } let cfg = &app_state.app_config; - let waker = Arc::new(AtomicWaker::new()); - let grace_stop_flag = - Self::stream_grace_period(app_state, &stream_details, grant_user_grace_period, user, fingerprint, Some(Arc::clone(&waker))); - let waker = if grace_stop_flag.is_some() { Some(waker) } else { None }; + let (grace_stop_flag, waker) = if grant_user_grace_period || (stream_details.has_grace_period() && stream_details.provider_name.is_some()) { + let waker = Arc::new(AtomicWaker::new()); + let flag = Self::stream_grace_period(app_state, &stream_details, grant_user_grace_period, user, fingerprint, Some(Arc::clone(&waker))); + let maybe_waker = flag.as_ref().map(|_| waker); + (flag, maybe_waker) + } else { + (Self::stream_grace_period(app_state, &stream_details, grant_user_grace_period, user, fingerprint, None), None) + }; + let custom_response = cfg.custom_stream_response.load(); let custom_video = custom_response.as_ref() .map_or((None, None, None), |c| @@ -105,6 +110,7 @@ impl ActiveClientStream { waker, connection_manager: Arc::clone(&app_state.connection_manager), fingerprint: Arc::new(fingerprint.clone()), + provider_stopped: false, } } @@ -174,7 +180,7 @@ impl ActiveClientStream { if updated { if let Some(flag) = reconnect_flag { - flag.notify(); + flag.notify(); } } @@ -187,20 +193,23 @@ impl ActiveClientStream { None } - fn stop_provider_stream(&mut self) { + fn stop_provider_stream(&mut self, unavailable: bool) { + self.provider_stopped = true; if self.provider_handle.is_some() { let mgr = Arc::clone(&self.connection_manager); let handle = self.provider_handle.take(); - if let Some(flag) = &self.send_custom_stream_flag { - flag.store(CHANNEL_UNAVAILABLE_STREAM, std::sync::atomic::Ordering::Release); + if unavailable { + if let Some(flag) = &self.send_custom_stream_flag { + flag.store(CHANNEL_UNAVAILABLE_STREAM, std::sync::atomic::Ordering::Release); + } } let con_man = Arc::clone(&self.connection_manager); let addr = self.fingerprint.addr; self.inner = futures::stream::empty::>().boxed(); tokio::spawn(async move { - con_man.update_stream_detail(&addr, CustomVideoStreamType::UserConnectionsExhausted).await; + con_man.update_stream_detail(&addr, if unavailable { CustomVideoStreamType::ChannelUnavailable } else { CustomVideoStreamType::UserConnectionsExhausted }).await; debug_if_enabled!("Provider stream stopped due to grace period or unavailable provider channel for {}", sanitize_sensitive_info(&addr.to_string())); mgr.release_provider_handle(handle).await; }); @@ -225,18 +234,19 @@ impl Stream for ActiveClientStream { match Pin::new(&mut self.inner).poll_next(cx) { Poll::Ready(Some(Err(e))) => { error!("Inner stream error: {e:?}"); - self.stop_provider_stream(); + self.stop_provider_stream(true); return Poll::Ready(Some(Err(e))); } Poll::Ready(None) => { - self.stop_provider_stream(); + self.stop_provider_stream(true); return Poll::Ready(None); } healthy => return healthy, } } - - self.stop_provider_stream(); + if !self.provider_stopped { + self.stop_provider_stream(false); + } let buffer_opt = match flag { USER_EXHAUSTED_STREAM => { @@ -252,10 +262,11 @@ impl Stream for ActiveClientStream { }; if let Some(buffer) = buffer_opt { - return Poll::Ready(Some(Ok(buffer.next_chunk()))); + Poll::Ready(Some(Ok(buffer.next_chunk()))) + } else { + // At this point it should be the empty stream + Pin::new(&mut self.inner).poll_next(cx) } - - Poll::Ready(None) } } diff --git a/backend/src/api/panel_api.rs b/backend/src/api/panel_api.rs index 014c68aa5..acc6b5987 100644 --- a/backend/src/api/panel_api.rs +++ b/backend/src/api/panel_api.rs @@ -1,7 +1,7 @@ use crate::api::model::AppState; use crate::model::{ConfigInput, is_input_expired, PanelApiQueryParam, PanelApiConfig}; -use crate::utils::debug_if_enabled; -use crate::utils::{format_sources_yaml_panel_api_query_params_flow_style, get_csv_file_path}; +use crate::utils::{debug_if_enabled, save_sources_config}; +use crate::utils::{get_csv_file_path}; use log::{debug, error, warn}; use serde_json::Value; use shared::error::{create_tuliprox_error_result, info_err, TuliproxError, TuliproxErrorKind}; @@ -289,7 +289,9 @@ fn collect_expired_accounts(input: &ConfigInput) -> Vec { out } +#[allow(clippy::too_many_arguments)] async fn patch_source_yml_add_alias( + app_state: &Arc, source_file_path: &Path, input_name: &str, alias_name: &str, @@ -370,11 +372,12 @@ async fn patch_source_yml_add_alias( } alias_seq.push(serde_yaml::Value::Mapping(alias_map)); - persist_source_config(source_file_path, &doc).await?; + persist_source_config(app_state, source_file_path, &doc).await?; Ok(()) } async fn patch_source_yml_update_exp_date( + app_state: &Arc, source_file_path: &Path, input_name: &str, account_name: &str, @@ -425,24 +428,20 @@ async fn patch_source_yml_update_exp_date( } } } - persist_source_config(source_file_path, &doc).await?; + persist_source_config(app_state, source_file_path, &doc).await?; return Ok(()); } } create_tuliprox_error_result!(TuliproxErrorKind::Info, "panel_api: could not find account '{account_name}' under input '{input_name}' in source.yml") } -async fn persist_source_config(source_file_path: &Path, doc: &serde_yaml::Value) -> Result<(), TuliproxError> { +async fn persist_source_config(app_state: &Arc, source_file_path: &Path, doc: &serde_yaml::Value) -> Result<(), TuliproxError> { - // TODO write backup!! Why not user ConfigFiles::save_sources_config - - let serialized = serde_yaml::to_string(&doc) - .map_err(|e| TuliproxError::new(TuliproxErrorKind::Info, format!("panel_api: failed to serialize source.yml: {e}")))?; - let serialized = format_sources_yaml_panel_api_query_params_flow_style(&serialized); - tokio::fs::write(source_file_path, serialized) - .await - .map_err(|e| TuliproxError::new(TuliproxErrorKind::Info, format!("panel_api: failed to write source.yml: {e}")))?; - Ok(()) + let paths = app_state.app_config.paths.load(); + let source_file = paths.sources_file_path.as_str(); + let config = app_state.app_config.config.load(); + let backup_dir = config.get_backup_dir(); + save_sources_config(source_file_path.to_str().unwrap_or(source_file), &backup_dir, &doc).await } async fn patch_batch_csv_append( @@ -608,7 +607,7 @@ async fn try_renew_expired_account( } } else { let _src_lock = app_state.app_config.file_locks.write_lock(sources_path).await; - if let Err(err) = patch_source_yml_update_exp_date(sources_path, &input.name, &acct.name, new_exp).await { + if let Err(err) = patch_source_yml_update_exp_date(app_state, sources_path, &input.name, &acct.name, new_exp).await { debug_if_enabled!("panel_api failed to persist renew exp_date to source.yml: {}", err); } } @@ -683,7 +682,7 @@ async fn try_create_new_account( } else { let _src_lock = app_state.app_config.file_locks.write_lock(sources_path).await; if let Err(err) = - patch_source_yml_add_alias(sources_path, &input.name, &alias_name, &base_url, &username, &password, exp_date).await + patch_source_yml_add_alias(app_state, sources_path, &input.name, &alias_name, &base_url, &username, &password, exp_date).await { warn!("panel_api failed to persist new alias to source.yml: {err}"); return false; @@ -811,7 +810,7 @@ pub(crate) async fn sync_panel_api_exp_dates_on_boot(app_state: &Arc) } } else { let _src_lock = app_state.app_config.file_locks.write_lock(&sources_path).await; - if let Err(err) = patch_source_yml_update_exp_date(&sources_path, &input.name, &acct.name, new_exp).await { + if let Err(err) = patch_source_yml_update_exp_date(app_state, &sources_path, &input.name, &acct.name, new_exp).await { debug_if_enabled!("panel_api boot sync failed to persist exp_date to source.yml: {}", err); continue; } @@ -830,11 +829,4 @@ pub(crate) async fn sync_panel_api_exp_dates_on_boot(app_state: &Arc) async fn reload_sources(app_state: &Arc) -> Result<(), TuliproxError> { ConfigFile::load_sources(app_state).await - // let paths = app_state.app_config.paths.load(); - // let sources_file = paths.sources_file_path.as_str(); - // let dto = read_sources_file(sources_file, true, true, None)?; - // let sources = crate::model::SourcesConfig::try_from(&dto)?; - // app_state.app_config.set_sources(sources)?; - // app_state.active_provider.update_config(&app_state.app_config).await; - // Ok(()) } \ No newline at end of file diff --git a/backend/src/model/config/panel_api.rs b/backend/src/model/config/panel_api.rs index a25a3a98f..feb8ce6b1 100644 --- a/backend/src/model/config/panel_api.rs +++ b/backend/src/model/config/panel_api.rs @@ -172,6 +172,9 @@ impl From<&PanelApiConfig> for PanelApiConfigDto { impl PanelApiConfig { pub fn prepare(&mut self) -> Result<(), TuliproxError> { + if !self.enabled { + return Ok(()); + } if self.url.trim().is_empty() { return create_tuliprox_error_result!(TuliproxErrorKind::Info, "panel_api: url is missing"); } diff --git a/backend/src/processing/processor/playlist.rs b/backend/src/processing/processor/playlist.rs index 01ca8ac10..6c2c72078 100644 --- a/backend/src/processing/processor/playlist.rs +++ b/backend/src/processing/processor/playlist.rs @@ -65,6 +65,7 @@ pub fn apply_favourites_to_playlist( _playlist: &mut [PlaylistGroup], _favourites_cfg: Option<&[ConfigFavourites]>, ) { + // TODO implement favourites // if let Some(favourites) = favourites_cfg { // let mut fav_groups: HashMap> = HashMap::new(); // diff --git a/backend/src/utils/file/config_reader.rs b/backend/src/utils/file/config_reader.rs index 5331b76fd..088664fc7 100644 --- a/backend/src/utils/file/config_reader.rs +++ b/backend/src/utils/file/config_reader.rs @@ -519,7 +519,10 @@ pub async fn save_main_config(file_path: &str, backup_dir: &str, config: &Config write_config_file(file_path, backup_dir, config, "config.yml", None).await } -pub async fn save_sources_config(file_path: &str, backup_dir: &str, config: &SourcesConfigDto) -> Result<(), TuliproxError> { +pub async fn save_sources_config(file_path: &str, backup_dir: &str, config: &T) -> Result<(), TuliproxError> +where + T: ?Sized + Serialize +{ write_config_file(file_path, backup_dir, config, "source.yml", Some(&format_sources_yaml_panel_api_query_params_flow_style)).await }