diff --git a/README.md b/README.md index e5d23bf01..d5e8efe48 100644 --- a/README.md +++ b/README.md @@ -826,6 +826,53 @@ Input alias definition for same provider with same content but different credent - name: test ``` +#### `panel_api` (optional) +If provider connections are exhausted, tuliprox can optionally call a provider panel API to: +- renew expired accounts first (based on `exp_date`) +- otherwise create a new alias account and persist it + +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: +- `api_key: auto` is replaced by `panel_api.api_key` +- in `renew_client`, `username: auto` / `password: auto` are replaced by the account being renewed + +`new_client` must not include a `user` query parameter; tuliprox expects the response to contain `username`/`password` or a `url` that contains them as query parameters. + +`panel_info` is used as a healthcheck and must return JSON with truthy `status` and `credits > 0`, e.g.: +```json +[{"status":"true","credits":"20","enabled":"1"}] +``` + +Example: +```yaml +- sources: +- inputs: + - type: xtream + name: my_provider + url: 'http://provider.net' + username: xyz + password: secret1 + panel_api: + url: 'https://panel.example.tld/api.php' + api_key: '${env:PANEL_API_KEY}' + query_parameter: + panel_info: + - { key: action, value: reseller_info } + - { key: api_key, value: auto } + new_client: + - { key: action, value: new } + - { key: type, value: m3u } + - { key: sub, value: '1' } + - { key: api_key, value: auto } + renew_client: + - { key: action, value: renew } + - { key: type, value: m3u } + - { key: username, value: auto } + - { key: password, value: auto } + - { key: sub, value: '1' } + - { key: api_key, value: auto } +``` + Input aliases can be defined as batches in csv files with `;` separator. There are 2 batch input types `xtream_batch` and `m3u_batch`. diff --git a/backend/src/api/api_utils.rs b/backend/src/api/api_utils.rs index 8c6c7a2ad..8d522593a 100644 --- a/backend/src/api/api_utils.rs +++ b/backend/src/api/api_utils.rs @@ -9,6 +9,7 @@ use crate::api::model::{ StreamError, ThrottledStream, UserApiRequest, }; use crate::api::model::{ProviderAllocation, ProviderConfig, ProviderStreamState, StreamDetails, StreamingStrategy}; +use crate::api::panel_api::try_provision_account_on_exhausted; use crate::model::{ConfigInput, ResourceRetryConfig}; use crate::model::{ConfigTarget, ProxyUserCredentials}; use crate::tools::lru_cache::LRUResourceCache; @@ -341,11 +342,20 @@ async fn resolve_streaming_strategy( force_provider: Option<&str>, ) -> StreamingStrategy { // allocate a provider connection - let provider_connection_handle = match force_provider { + let mut provider_connection_handle = match force_provider { Some(provider) => app_state.active_provider.force_exact_acquire_connection(provider, &fingerprint.addr).await, None => app_state.active_provider.acquire_connection(&input.name, &fingerprint.addr).await, }; + if provider_connection_handle.is_none() && force_provider.is_none() { + if try_provision_account_on_exhausted(app_state, input).await { + provider_connection_handle = app_state + .active_provider + .acquire_connection(&input.name, &fingerprint.addr) + .await; + } + } + let stream_response_params = if let Some(allocation) = provider_connection_handle.as_ref().map(|ph| &ph.allocation) { match allocation { diff --git a/backend/src/api/mod.rs b/backend/src/api/mod.rs index 3789ade60..1741f4570 100644 --- a/backend/src/api/mod.rs +++ b/backend/src/api/mod.rs @@ -1,5 +1,6 @@ pub mod model; pub mod api_utils; +mod panel_api; mod scheduler; mod endpoints; pub mod main_api; @@ -7,4 +8,4 @@ mod config_watch; mod serve; pub(crate) mod hdhomerun_ssdp; pub(crate) mod hdhomerun_proprietary; -mod sys_usage; \ No newline at end of file +mod sys_usage; diff --git a/backend/src/api/model/provider_lineup_manager.rs b/backend/src/api/model/provider_lineup_manager.rs index b25455cce..e29093d2a 100644 --- a/backend/src/api/model/provider_lineup_manager.rs +++ b/backend/src/api/model/provider_lineup_manager.rs @@ -803,28 +803,29 @@ mod tests { } // Helper function to create a ConfigInput instance - fn create_config_input(id: u16, name: &str, priority: i16, max_connections: u16) -> ConfigInput { - ConfigInput { - id, - name: name.to_string(), - url: "http://example.com".to_string(), - epg: Option::default(), - username: None, - password: None, - persist: None, - enabled: true, - input_type: InputType::Xtream, // You can use a default value here - max_connections, - priority, - aliases: None, - headers: HashMap::default(), - options: None, - method: InputFetchMethod::default(), - staged: None, - exp_date: None, - t_batch_url: None, - } - } + fn create_config_input(id: u16, name: &str, priority: i16, max_connections: u16) -> ConfigInput { + ConfigInput { + id, + name: name.to_string(), + url: "http://example.com".to_string(), + epg: Option::default(), + username: None, + password: None, + persist: None, + enabled: true, + input_type: InputType::Xtream, // You can use a default value here + max_connections, + priority, + aliases: None, + headers: HashMap::default(), + options: None, + method: InputFetchMethod::default(), + staged: None, + exp_date: None, + t_batch_url: None, + panel_api: None, + } + } // Helper function to create a ConfigInputAlias instance fn create_config_input_alias(id: u16, url: &str, priority: i16, max_connections: u16) -> ConfigInputAlias { diff --git a/backend/src/api/panel_api.rs b/backend/src/api/panel_api.rs new file mode 100644 index 000000000..858c4a040 --- /dev/null +++ b/backend/src/api/panel_api.rs @@ -0,0 +1,676 @@ +use crate::api::model::AppState; +use crate::model::{ConfigInput, is_input_expired}; +use crate::utils::debug_if_enabled; +use crate::utils::{get_csv_file_path, read_sources_file}; +use log::{debug, error, warn}; +use serde_json::Value; +use shared::error::{create_tuliprox_error_result, info_err, TuliproxError, TuliproxErrorKind}; +use shared::model::{InputType, PanelApiConfigDto, PanelApiQueryParamDto}; +use shared::utils::{get_credentials_from_url, sanitize_sensitive_info, trim_last_slash}; +use std::collections::HashSet; +use std::path::{Path, PathBuf}; +use url::Url; + +#[derive(Debug, Clone)] +struct AccountCredentials { + name: String, + username: String, + password: String, + exp_date: Option, +} + +fn parse_boolish(value: &Value) -> bool { + match value { + Value::Bool(b) => *b, + Value::Number(n) => n.as_i64().unwrap_or(0) != 0, + Value::String(s) => matches!(s.trim().to_lowercase().as_str(), "true" | "1" | "yes" | "y" | "ok"), + _ => false, + } +} + +fn parse_i64ish(value: &Value) -> Option { + match value { + Value::Number(n) => n.as_i64(), + Value::String(s) => s.trim().parse::().ok(), + _ => None, + } +} + +fn first_json_object(value: &Value) -> Option<&serde_json::Map> { + match value { + Value::Array(arr) => arr.first().and_then(|v| v.as_object()), + Value::Object(obj) => Some(obj), + _ => None, + } +} + +fn extract_username_password_from_json(obj: &serde_json::Map) -> Option<(String, String)> { + let username = obj.get("username").and_then(|v| v.as_str()).map(str::trim).filter(|s| !s.is_empty()); + let password = obj.get("password").and_then(|v| v.as_str()).map(str::trim).filter(|s| !s.is_empty()); + match (username, password) { + (Some(u), Some(p)) => Some((u.to_string(), p.to_string())), + _ => None, + } +} + +fn extract_username_password_from_url(url_str: &str) -> Option<(String, String)> { + Url::parse(url_str).ok().and_then(|url| { + let (u, p) = get_credentials_from_url(&url); + match (u, p) { + (Some(u), Some(p)) if !u.trim().is_empty() && !p.trim().is_empty() => Some((u, p)), + _ => None, + } + }) +} + +fn extract_base_url(url_str: &str) -> Option { + Url::parse(url_str).ok().map(|u| u.origin().ascii_serialization()) +} + +fn validate_type_is_m3u(params: &[PanelApiQueryParamDto]) -> Result<(), TuliproxError> { + let typ = params + .iter() + .find(|p| p.key.trim().eq_ignore_ascii_case("type")) + .map(|p| p.value.trim().to_string()); + match typ { + Some(v) if v.eq_ignore_ascii_case("m3u") => Ok(()), + Some(v) => create_tuliprox_error_result!(TuliproxErrorKind::Info, "panel_api: unsupported type={v}, only m3u is supported"), + None => create_tuliprox_error_result!(TuliproxErrorKind::Info, "panel_api: missing required query param 'type=m3u'"), + } +} + +fn validate_new_client_params(params: &[PanelApiQueryParamDto]) -> Result<(), TuliproxError> { + validate_type_is_m3u(params)?; + if params.iter().any(|p| p.key.trim().eq_ignore_ascii_case("user")) { + return create_tuliprox_error_result!(TuliproxErrorKind::Info, "panel_api: new_client must not contain query param 'user'"); + } + Ok(()) +} + +fn validate_renew_client_params(params: &[PanelApiQueryParamDto]) -> Result<(), TuliproxError> { + validate_type_is_m3u(params)?; + let keys: HashSet = params.iter().map(|p| p.key.trim().to_lowercase()).collect(); + if !(keys.contains("username") && keys.contains("password")) { + return create_tuliprox_error_result!(TuliproxErrorKind::Info, "panel_api: renew_client must contain query params 'username' and 'password' (use value 'auto')"); + } + Ok(()) +} + +fn resolve_query_params( + params: &[PanelApiQueryParamDto], + api_key: Option<&str>, + creds: Option<(&str, &str)>, +) -> Result, TuliproxError> { + let mut out = Vec::with_capacity(params.len()); + for p in params { + let key = p.key.trim(); + if key.is_empty() { + continue; + } + let mut value = p.value.trim().to_string(); + if value.eq_ignore_ascii_case("auto") { + if key.eq_ignore_ascii_case("api_key") { + let Some(k) = api_key.filter(|s| !s.trim().is_empty()) else { + return create_tuliprox_error_result!(TuliproxErrorKind::Info, "panel_api: query param {key} uses 'auto' but panel_api.api_key is missing"); + }; + value = k.to_string(); + } else if key.eq_ignore_ascii_case("username") { + let Some((u, _)) = creds else { + return create_tuliprox_error_result!(TuliproxErrorKind::Info, "panel_api: query param {key} uses 'auto' but no account username is available"); + }; + value = u.to_string(); + } else if key.eq_ignore_ascii_case("password") { + let Some((_, pw)) = creds else { + return create_tuliprox_error_result!(TuliproxErrorKind::Info, "panel_api: query param {key} uses 'auto' but no account password is available"); + }; + value = pw.to_string(); + } + } + out.push((key.to_string(), value)); + } + Ok(out) +} + +fn build_panel_url(base_url: &str, query_params: &[(String, String)]) -> Result { + let mut url = Url::parse(base_url).map_err(|e| info_err!(format!("panel_api: invalid url {base_url}: {e}")))?; + { + let mut pairs = url.query_pairs_mut(); + for (k, v) in query_params { + pairs.append_pair(k, v); + } + } + Ok(url) +} + +async fn panel_get_json(app_state: &AppState, url: Url) -> Result { + let client = app_state.http_client.load(); + let sanitized = sanitize_sensitive_info(url.as_str()); + debug_if_enabled!("panel_api request {}", sanitized); + let resp = client + .get(url) + .send() + .await + .map_err(|e| TuliproxError::new(TuliproxErrorKind::Info, format!("panel_api request failed: {e}")))?; + let status = resp.status(); + let body = resp + .text() + .await + .map_err(|e| TuliproxError::new(TuliproxErrorKind::Info, format!("panel_api read response failed: {e}")))?; + let json: Value = serde_json::from_str(&body) + .map_err(|e| TuliproxError::new(TuliproxErrorKind::Info, format!("panel_api invalid json (http {status}): {e}")))?; + Ok(json) +} + +async fn panel_healthcheck(app_state: &AppState, cfg: &PanelApiConfigDto) -> Result<(), TuliproxError> { + if cfg.query_parameter.panel_info.is_empty() { + return create_tuliprox_error_result!(TuliproxErrorKind::Info, "panel_api: query_parameter.panel_info is missing/empty"); + } + let params = resolve_query_params(&cfg.query_parameter.panel_info, cfg.api_key.as_deref(), None)?; + let url = build_panel_url(cfg.url.as_str(), ¶ms)?; + let json = panel_get_json(app_state, url).await?; + let Some(obj) = first_json_object(&json) else { + return create_tuliprox_error_result!(TuliproxErrorKind::Info, "panel_api: panel_info response is not a JSON object/array"); + }; + let status_ok = obj.get("status").map(parse_boolish).unwrap_or(false); + let credits = obj.get("credits").and_then(parse_i64ish).unwrap_or(0); + if !status_ok || credits <= 0 { + return create_tuliprox_error_result!(TuliproxErrorKind::Info, "panel_api: healthcheck failed (status={status_ok}, credits={credits})"); + } + Ok(()) +} + +async fn panel_new_client(app_state: &AppState, cfg: &PanelApiConfigDto) -> Result<(String, String, Option), TuliproxError> { + validate_new_client_params(&cfg.query_parameter.new_client)?; + let params = resolve_query_params(&cfg.query_parameter.new_client, cfg.api_key.as_deref(), None)?; + let url = build_panel_url(cfg.url.as_str(), ¶ms)?; + let json = panel_get_json(app_state, url).await?; + let Some(obj) = first_json_object(&json) else { + return create_tuliprox_error_result!(TuliproxErrorKind::Info, "panel_api: new_client response is not a JSON object/array"); + }; + let status_ok = obj.get("status").map(parse_boolish).unwrap_or(false); + if !status_ok { + return create_tuliprox_error_result!(TuliproxErrorKind::Info, "panel_api: new_client status=false"); + } + if let Some((u, p)) = extract_username_password_from_json(obj) { + return Ok((u, p, None)); + } + if let Some(url_str) = obj.get("url").and_then(|v| v.as_str()) { + if let Some((u, p)) = extract_username_password_from_url(url_str) { + let base = extract_base_url(url_str); + return Ok((u, p, base)); + } + } + create_tuliprox_error_result!(TuliproxErrorKind::Info, "panel_api: new_client response missing username/password (and no parsable url)") +} + +async fn panel_renew_client(app_state: &AppState, cfg: &PanelApiConfigDto, username: &str, password: &str) -> Result<(), TuliproxError> { + validate_renew_client_params(&cfg.query_parameter.renew_client)?; + let params = resolve_query_params( + &cfg.query_parameter.renew_client, + cfg.api_key.as_deref(), + Some((username, password)), + )?; + let url = build_panel_url(cfg.url.as_str(), ¶ms)?; + let json = panel_get_json(app_state, url).await?; + let Some(obj) = first_json_object(&json) else { + return create_tuliprox_error_result!(TuliproxErrorKind::Info, "panel_api: renew_client response is not a JSON object/array"); + }; + let status_ok = obj.get("status").map(parse_boolish).unwrap_or(false); + if !status_ok { + return create_tuliprox_error_result!(TuliproxErrorKind::Info, "panel_api: renew_client status=false"); + } + Ok(()) +} + +fn derive_exp_date_from_params(params: &[PanelApiQueryParamDto]) -> Option { + use chrono::{Duration, Months, Utc}; + let now = Utc::now(); + let raw = params + .iter() + .find(|p| { + let k = p.key.trim(); + k.eq_ignore_ascii_case("sub") || k.eq_ignore_ascii_case("package") + }) + .map(|p| p.value.trim().to_string())?; + let months = raw.parse::().ok()?; + if months <= 0 { + return Some((now + Duration::days(1)).timestamp()); + } + if months == 99 { + return Some((now + Duration::days(1)).timestamp()); + } + let months_u32: u32 = months.try_into().ok()?; + now.checked_add_months(Months::new(months_u32)).map(|dt| dt.timestamp()) +} + +fn extract_account_creds_from_input(input: &ConfigInput) -> Option<(String, String)> { + if let (Some(u), Some(p)) = (input.username.as_deref(), input.password.as_deref()) { + if !u.trim().is_empty() && !p.trim().is_empty() { + return Some((u.to_string(), p.to_string())); + } + } + Url::parse(input.url.as_str()).ok().and_then(|u| { + let (uu, pp) = get_credentials_from_url(&u); + match (uu, pp) { + (Some(uu), Some(pp)) if !uu.trim().is_empty() && !pp.trim().is_empty() => Some((uu, pp)), + _ => None, + } + }) +} + +fn collect_expired_accounts(input: &ConfigInput) -> Vec { + let mut out = Vec::new(); + if is_input_expired(input.exp_date) { + if let Some((u, p)) = extract_account_creds_from_input(input) { + out.push(AccountCredentials { + name: input.name.clone(), + username: u, + password: p, + exp_date: input.exp_date, + }); + } + } + if let Some(aliases) = input.aliases.as_ref() { + for a in aliases { + if is_input_expired(a.exp_date) { + if let (Some(u), Some(p)) = (a.username.as_deref(), a.password.as_deref()) { + if !u.trim().is_empty() && !p.trim().is_empty() { + out.push(AccountCredentials { + name: a.name.clone(), + username: u.to_string(), + password: p.to_string(), + exp_date: a.exp_date, + }); + } + } + } + } + } + out.sort_by_key(|a| a.exp_date.unwrap_or(i64::MAX)); + out +} + +async fn patch_source_yml_add_alias( + source_file_path: &Path, + input_name: &str, + alias_name: &str, + base_url: &str, + username: &str, + password: &str, + exp_date: Option, +) -> Result<(), TuliproxError> { + let raw = tokio::fs::read_to_string(source_file_path) + .await + .map_err(|e| TuliproxError::new(TuliproxErrorKind::Info, format!("panel_api: failed to read source file: {e}")))?; + let mut doc: serde_yaml::Value = serde_yaml::from_str(&raw) + .map_err(|e| TuliproxError::new(TuliproxErrorKind::Info, format!("panel_api: failed to parse source file yaml: {e}")))?; + + let Some(root) = doc.as_mapping_mut() else { + return create_tuliprox_error_result!(TuliproxErrorKind::Info, "panel_api: source.yml root is not a mapping"); + }; + let sources = root.get_mut(serde_yaml::Value::String("sources".to_string())).and_then(|v| v.as_sequence_mut()); + let Some(sources) = sources else { + return create_tuliprox_error_result!(TuliproxErrorKind::Info, "panel_api: source.yml missing 'sources' list"); + }; + + let mut found_input = None; + for src in sources.iter_mut() { + let Some(src_map) = src.as_mapping_mut() else { continue; }; + let Some(inputs) = src_map.get_mut(serde_yaml::Value::String("inputs".to_string())).and_then(|v| v.as_sequence_mut()) else { continue; }; + for inp in inputs.iter_mut() { + let Some(inp_map) = inp.as_mapping_mut() else { continue; }; + let name = inp_map.get(serde_yaml::Value::String("name".to_string())).and_then(|v| v.as_str()); + if name == Some(input_name) { + found_input = Some(inp_map); + break; + } + } + if found_input.is_some() { + break; + } + } + let Some(inp_map) = found_input else { + return create_tuliprox_error_result!(TuliproxErrorKind::Info, "panel_api: could not find input '{input_name}' in source.yml"); + }; + + let aliases_key = serde_yaml::Value::String("aliases".to_string()); + if !inp_map.contains_key(&aliases_key) { + inp_map.insert(aliases_key.clone(), serde_yaml::Value::Sequence(vec![])); + } + let Some(alias_seq) = inp_map.get_mut(&aliases_key).and_then(|v| v.as_sequence_mut()) else { + return create_tuliprox_error_result!(TuliproxErrorKind::Info, "panel_api: input.aliases is not a list in source.yml"); + }; + + let mut alias_map = serde_yaml::Mapping::new(); + alias_map.insert(serde_yaml::Value::String("name".to_string()), serde_yaml::Value::String(alias_name.to_string())); + alias_map.insert(serde_yaml::Value::String("url".to_string()), serde_yaml::Value::String(base_url.to_string())); + alias_map.insert(serde_yaml::Value::String("username".to_string()), serde_yaml::Value::String(username.to_string())); + alias_map.insert(serde_yaml::Value::String("password".to_string()), serde_yaml::Value::String(password.to_string())); + alias_map.insert(serde_yaml::Value::String("max_connections".to_string()), serde_yaml::Value::Number(1.into())); + if let Some(ts) = exp_date { + alias_map.insert(serde_yaml::Value::String("exp_date".to_string()), serde_yaml::Value::Number(ts.into())); + } + alias_seq.push(serde_yaml::Value::Mapping(alias_map)); + + let serialized = serde_yaml::to_string(&doc) + .map_err(|e| TuliproxError::new(TuliproxErrorKind::Info, format!("panel_api: failed to serialize source.yml: {e}")))?; + 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(()) +} + +async fn patch_source_yml_update_exp_date( + source_file_path: &Path, + input_name: &str, + account_name: &str, + exp_date: i64, +) -> Result<(), TuliproxError> { + let raw = tokio::fs::read_to_string(source_file_path) + .await + .map_err(|e| TuliproxError::new(TuliproxErrorKind::Info, format!("panel_api: failed to read source file: {e}")))?; + let mut doc: serde_yaml::Value = serde_yaml::from_str(&raw) + .map_err(|e| TuliproxError::new(TuliproxErrorKind::Info, format!("panel_api: failed to parse source file yaml: {e}")))?; + let Some(root) = doc.as_mapping_mut() else { + return create_tuliprox_error_result!(TuliproxErrorKind::Info, "panel_api: source.yml root is not a mapping"); + }; + let sources = root.get_mut(serde_yaml::Value::String("sources".to_string())).and_then(|v| v.as_sequence_mut()); + let Some(sources) = sources else { + return create_tuliprox_error_result!(TuliproxErrorKind::Info, "panel_api: source.yml missing 'sources' list"); + }; + for src in sources.iter_mut() { + let Some(src_map) = src.as_mapping_mut() else { continue; }; + let Some(inputs) = src_map.get_mut(serde_yaml::Value::String("inputs".to_string())).and_then(|v| v.as_sequence_mut()) else { continue; }; + for inp in inputs.iter_mut() { + let Some(inp_map) = inp.as_mapping_mut() else { continue; }; + let name = inp_map.get(serde_yaml::Value::String("name".to_string())).and_then(|v| v.as_str()); + if name != Some(input_name) { + continue; + } + if account_name == input_name { + inp_map.insert(serde_yaml::Value::String("exp_date".to_string()), serde_yaml::Value::Number(exp_date.into())); + inp_map.insert(serde_yaml::Value::String("enabled".to_string()), serde_yaml::Value::Bool(true)); + } else if let Some(aliases) = inp_map.get_mut(serde_yaml::Value::String("aliases".to_string())).and_then(|v| v.as_sequence_mut()) { + for a in aliases.iter_mut() { + let Some(a_map) = a.as_mapping_mut() else { continue; }; + let a_name = a_map.get(serde_yaml::Value::String("name".to_string())).and_then(|v| v.as_str()); + if a_name == Some(account_name) { + a_map.insert(serde_yaml::Value::String("exp_date".to_string()), serde_yaml::Value::Number(exp_date.into())); + } + } + } + let serialized = serde_yaml::to_string(&doc) + .map_err(|e| TuliproxError::new(TuliproxErrorKind::Info, format!("panel_api: failed to serialize source.yml: {e}")))?; + tokio::fs::write(source_file_path, serialized) + .await + .map_err(|e| TuliproxError::new(TuliproxErrorKind::Info, format!("panel_api: failed to write source.yml: {e}")))?; + 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 patch_batch_csv_append( + csv_path: &Path, + batch_type: InputType, + alias_name: &str, + base_url: &str, + username: &str, + password: &str, + exp_date: Option, +) -> Result<(), TuliproxError> { + let raw = tokio::fs::read_to_string(csv_path).await.unwrap_or_default(); + let mut lines: Vec = raw.lines().map(|l| l.to_string()).collect(); + let header_line_idx = lines.iter().position(|l| l.trim_start().starts_with('#')); + let header = header_line_idx + .and_then(|idx| lines.get(idx).map(|s| s.trim_start_matches('#').trim().to_string())) + .unwrap_or_else(|| match batch_type { + InputType::XtreamBatch => "name;username;password;url;max_connections;priority;exp_date".to_string(), + _ => "url;max_connections;priority;name;username;password;exp_date".to_string(), + }); + let cols: Vec = header.split(';').map(|s| s.trim().to_lowercase()).collect(); + if header_line_idx.is_none() { + lines.insert(0, format!("#{header}")); + } + let mut row: Vec = vec![String::new(); cols.len()]; + for (i, c) in cols.iter().enumerate() { + row[i] = match c.as_str() { + "name" => alias_name.to_string(), + "username" => username.to_string(), + "password" => password.to_string(), + "url" => { + if batch_type == InputType::M3uBatch { + format!( + "{}/get.php?username={}&password={}&type=m3u_plus", + trim_last_slash(base_url), + username, + password + ) + } else { + base_url.to_string() + } + } + "max_connections" => "1".to_string(), + "priority" => "0".to_string(), + "exp_date" => exp_date.map_or(String::new(), |ts| ts.to_string()), + _ => String::new(), + }; + } + lines.push(row.join(";")); + tokio::fs::write(csv_path, lines.join("\n") + "\n") + .await + .map_err(|e| TuliproxError::new(TuliproxErrorKind::Info, format!("panel_api: failed to write csv: {e}")))?; + Ok(()) +} + +async fn patch_batch_csv_update_exp_date( + csv_path: &Path, + account_name: &str, + username: &str, + password: &str, + exp_date: i64, +) -> Result<(), TuliproxError> { + let raw = tokio::fs::read_to_string(csv_path) + .await + .map_err(|e| TuliproxError::new(TuliproxErrorKind::Info, format!("panel_api: failed to read csv: {e}")))?; + let mut lines: Vec = raw.lines().map(|l| l.to_string()).collect(); + let header_line_idx = lines.iter().position(|l| l.trim_start().starts_with('#')); + let Some(header_idx) = header_line_idx else { + return create_tuliprox_error_result!(TuliproxErrorKind::Info, "panel_api: csv missing header line"); + }; + let header = lines[header_idx].trim_start_matches('#').trim(); + let cols: Vec = header.split(';').map(|s| s.trim().to_lowercase()).collect(); + let exp_idx = cols.iter().position(|c| c == "exp_date"); + let name_idx = cols.iter().position(|c| c == "name"); + let user_idx = cols.iter().position(|c| c == "username"); + let pass_idx = cols.iter().position(|c| c == "password"); + let url_idx = cols.iter().position(|c| c == "url"); + let Some(exp_idx) = exp_idx else { + debug!("panel_api: csv has no exp_date column; skipping exp_date persistence"); + return Ok(()); + }; + + for i in (header_idx + 1)..lines.len() { + let line = lines[i].trim(); + if line.is_empty() || line.starts_with('#') { + continue; + } + let mut fields: Vec = line.split(';').map(|s| s.to_string()).collect(); + fields.resize(cols.len(), String::new()); + + let mut matches = false; + if let Some(n_idx) = name_idx { + if fields.get(n_idx).map(|s| s.trim()) == Some(account_name) { + matches = true; + } + } + if !matches { + if let (Some(u_idx), Some(p_idx)) = (user_idx, pass_idx) { + matches = fields.get(u_idx).map(|s| s.trim()) == Some(username) + && fields.get(p_idx).map(|s| s.trim()) == Some(password); + } else if let Some(u_idx) = url_idx { + if let Some(url_str) = fields.get(u_idx) { + if let Some((u, p)) = extract_username_password_from_url(url_str) { + matches = u == username && p == password; + } + } + } + } + if matches { + fields[exp_idx] = exp_date.to_string(); + lines[i] = fields.join(";"); + tokio::fs::write(csv_path, lines.join("\n") + "\n") + .await + .map_err(|e| TuliproxError::new(TuliproxErrorKind::Info, format!("panel_api: failed to write csv: {e}")))?; + return Ok(()); + } + } + warn!("panel_api: could not find batch csv row for account {}", account_name); + Ok(()) +} + +fn derive_unique_alias_name(existing: &[String], input_name: &str, username: &str) -> String { + let base = format!("{input_name}-{username}"); + if !existing.contains(&base) { + return base; + } + for i in 2..1000 { + let cand = format!("{base}-{i}"); + if !existing.contains(&cand) { + return cand; + } + } + base +} + +pub async fn try_provision_account_on_exhausted(app_state: &AppState, input: &ConfigInput) -> bool { + let Some(panel_cfg) = input.panel_api.as_ref() else { + return false; + }; + if panel_cfg.url.trim().is_empty() { + return false; + } + + let _input_lock = app_state + .app_config + .file_locks + .write_lock_str(format!("panel_api:{}", input.name).as_str()) + .await; + + if let Err(err) = panel_healthcheck(app_state, panel_cfg).await { + debug_if_enabled!("panel_api healthcheck failed: {}", sanitize_sensitive_info(err.to_string().as_str())); + return false; + } + + let is_batch = input.t_batch_url.as_ref().is_some_and(|u| !u.trim().is_empty()); + let sources_file_path = app_state.app_config.paths.load().sources_file_path.clone(); + let sources_path = PathBuf::from(&sources_file_path); + + // prefer renew of expired accounts + let expired = collect_expired_accounts(input); + if !expired.is_empty() { + for acct in &expired { + match panel_renew_client(app_state, panel_cfg, acct.username.as_str(), acct.password.as_str()).await { + Ok(()) => { + let new_exp = derive_exp_date_from_params(&panel_cfg.query_parameter.renew_client) + .or_else(|| derive_exp_date_from_params(&panel_cfg.query_parameter.new_client)); + if let Some(new_exp) = new_exp { + if is_batch { + let batch_url = input.t_batch_url.as_deref().unwrap_or_default(); + if let Ok(csv_path) = get_csv_file_path(batch_url) { + let _csv_lock = app_state.app_config.file_locks.write_lock(&csv_path).await; + if let Err(err) = patch_batch_csv_update_exp_date(&csv_path, &acct.name, &acct.username, &acct.password, new_exp).await { + debug_if_enabled!("panel_api failed to persist renew exp_date to csv: {}", err); + } + } + } 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 { + debug_if_enabled!("panel_api failed to persist renew exp_date to source.yml: {}", err); + } + } + } + + // reload sources + provider lineup + if let Err(err) = reload_sources(app_state).await { + debug_if_enabled!("panel_api reload sources failed: {}", err); + } + return true; + } + Err(err) => { + debug_if_enabled!( + "panel_api renew failed for {}: {}", + sanitize_sensitive_info(&acct.name), + sanitize_sensitive_info(err.to_string().as_str()) + ); + continue; + } + } + } + } + + // otherwise create new client and append as alias / batch row + match panel_new_client(app_state, panel_cfg).await { + Ok((username, password, base_url_from_resp)) => { + let base_url = base_url_from_resp.unwrap_or_else(|| input.url.clone()); + let base_url = extract_base_url(base_url.as_str()).unwrap_or_else(|| base_url.clone()); + + let mut existing_names: Vec = vec![input.name.clone()]; + if let Some(aliases) = input.aliases.as_ref() { + existing_names.extend(aliases.iter().map(|a| a.name.clone())); + } + let alias_name = derive_unique_alias_name(&existing_names, &input.name, &username); + + let exp_date = derive_exp_date_from_params(&panel_cfg.query_parameter.new_client) + .or_else(|| derive_exp_date_from_params(&panel_cfg.query_parameter.renew_client)); + + if is_batch { + let batch_url = input.t_batch_url.as_deref().unwrap_or_default(); + match get_csv_file_path(batch_url) { + Ok(csv_path) => { + let batch_type = if input.input_type == InputType::Xtream { + InputType::XtreamBatch + } else { + InputType::M3uBatch + }; + let _csv_lock = app_state.app_config.file_locks.write_lock(&csv_path).await; + if let Err(err) = patch_batch_csv_append(&csv_path, batch_type, &alias_name, &base_url, &username, &password, exp_date).await { + warn!("panel_api failed to append new account to csv: {}", err); + return false; + } + } + Err(err) => { + warn!("panel_api cannot resolve batch csv path {}: {}", sanitize_sensitive_info(batch_url), err); + return false; + } + } + } 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 { + warn!("panel_api failed to persist new alias to source.yml: {}", err); + return false; + } + } + + if let Err(err) = reload_sources(app_state).await { + error!("panel_api reload sources failed: {err}"); + return false; + } + true + } + Err(err) => { + debug_if_enabled!("panel_api new_client failed: {}", sanitize_sensitive_info(err.to_string().as_str())); + false + } + } +} + +async fn reload_sources(app_state: &AppState) -> Result<(), TuliproxError> { + 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(()) +} diff --git a/backend/src/model/config/input.rs b/backend/src/model/config/input.rs index e2d8a1040..8425f8170 100644 --- a/backend/src/model/config/input.rs +++ b/backend/src/model/config/input.rs @@ -3,7 +3,7 @@ use crate::utils::get_csv_file_path; use chrono::Utc; use log::warn; use shared::error::TuliproxError; -use shared::model::{ConfigInputAliasDto, ConfigInputDto, ConfigInputOptionsDto, InputFetchMethod, InputType, StagedInputDto}; +use shared::model::{ConfigInputAliasDto, ConfigInputDto, ConfigInputOptionsDto, InputFetchMethod, InputType, PanelApiConfigDto, StagedInputDto}; use shared::utils::{get_base_url_from_str, get_credentials_from_url}; use shared::{check_input_connections, info_err, write_if_some}; use shared::check_input_credentials; @@ -142,6 +142,7 @@ pub struct ConfigInput { pub staged: Option, pub exp_date: Option, pub t_batch_url: Option, + pub panel_api: Option, } impl ConfigInput { @@ -242,6 +243,7 @@ impl ConfigInput { staged: None, exp_date: None, t_batch_url: None, + panel_api: self.panel_api.clone(), } } } @@ -268,6 +270,7 @@ impl From<&ConfigInputDto> for ConfigInput { exp_date: dto.exp_date, staged: dto.staged.as_ref().map(StagedInput::from), t_batch_url: None, + panel_api: dto.panel_api.clone(), } } } @@ -305,4 +308,4 @@ pub fn is_input_expired(exp_date: Option) -> bool { } None => false, } -} \ No newline at end of file +} diff --git a/shared/src/model/config/input.rs b/shared/src/model/config/input.rs index e2347715c..b54690243 100644 --- a/shared/src/model/config/input.rs +++ b/shared/src/model/config/input.rs @@ -1,6 +1,7 @@ use crate::error::{TuliproxError, TuliproxErrorKind}; use crate::model::{EpgConfigDto}; use crate::utils::{default_as_true, get_credentials_from_url_str, get_trimmed_string, sanitize_sensitive_info, trim_last_slash, deserialize_timestamp}; +use super::PanelApiConfigDto; use crate::{check_input_credentials, check_input_connections, create_tuliprox_error_result, handle_tuliprox_error_result_list, info_err}; use enum_iterator::Sequence; use std::collections::{HashMap, HashSet}; @@ -300,6 +301,8 @@ pub struct ConfigInputDto { pub staged: Option, #[serde(default, deserialize_with = "deserialize_timestamp", skip_serializing_if = "Option::is_none")] pub exp_date: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub panel_api: Option, } impl Default for ConfigInputDto { @@ -322,6 +325,7 @@ impl Default for ConfigInputDto { method: InputFetchMethod::default(), staged: None, exp_date: None, + panel_api: None, } } } diff --git a/shared/src/model/config/mod.rs b/shared/src/model/config/mod.rs index 0b10eda52..ced7ad4d6 100644 --- a/shared/src/model/config/mod.rs +++ b/shared/src/model/config/mod.rs @@ -12,6 +12,7 @@ mod video_download; mod schedule; mod log; mod input; +mod panel_api; mod stream; mod epg; mod reverse_proxy; @@ -53,6 +54,7 @@ pub use video_download::*; pub use schedule::*; pub use log::*; pub use input::*; +pub use panel_api::*; pub use stream::*; pub use epg_smart_match::*; pub use epg::*; @@ -69,4 +71,4 @@ pub use app_config::*; pub use config_type::*; pub use playlist_update_state::*; pub use favourites::*; -pub use crate::apply_batch_aliases; \ No newline at end of file +pub use crate::apply_batch_aliases; diff --git a/shared/src/model/config/panel_api.rs b/shared/src/model/config/panel_api.rs new file mode 100644 index 000000000..934b05d74 --- /dev/null +++ b/shared/src/model/config/panel_api.rs @@ -0,0 +1,28 @@ +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, Default, PartialEq)] +#[serde(deny_unknown_fields)] +pub struct PanelApiQueryParamDto { + pub key: String, + pub value: String, +} + +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, Default, PartialEq)] +#[serde(deny_unknown_fields)] +pub struct PanelApiQueryParametersDto { + #[serde(default, skip_serializing_if = "Vec::is_empty")] + pub panel_info: Vec, + #[serde(default, skip_serializing_if = "Vec::is_empty")] + pub new_client: Vec, + #[serde(default, skip_serializing_if = "Vec::is_empty")] + pub renew_client: Vec, +} + +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, Default, PartialEq)] +#[serde(deny_unknown_fields)] +pub struct PanelApiConfigDto { + pub url: String, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub api_key: Option, + #[serde(default)] + pub query_parameter: PanelApiQueryParametersDto, +} +