From 32a849d1e94b5984f885dbd643abea0f4be2c803 Mon Sep 17 00:00:00 2001 From: euzu Date: Fri, 16 May 2025 10:23:38 +0200 Subject: [PATCH] Check xc user account status during update --- src/api/endpoints/api_playlist_utils.rs | 2 +- src/model/api_proxy.rs | 1 + src/processing/processor/playlist.rs | 2 +- src/utils/json_utils.rs | 8 +++ src/utils/network/xtream.rs | 85 +++++++++++++++++++++---- 5 files changed, 84 insertions(+), 14 deletions(-) diff --git a/src/api/endpoints/api_playlist_utils.rs b/src/api/endpoints/api_playlist_utils.rs index 00460fd37..20744c540 100644 --- a/src/api/endpoints/api_playlist_utils.rs +++ b/src/api/endpoints/api_playlist_utils.rs @@ -171,7 +171,7 @@ pub(in crate::api::endpoints) async fn get_playlist(client: Arc let (result, errors) = match input.input_type { InputType::M3u | InputType::M3uBatch => m3u::get_m3u_playlist(client, cfg, input, &cfg.working_dir).await, - InputType::Xtream | InputType::XtreamBatch => xtream::get_xtream_playlist(client, input, &cfg.working_dir).await, + InputType::Xtream | InputType::XtreamBatch => xtream::get_xtream_playlist(cfg, client, input, &cfg.working_dir).await, }; if result.is_empty() { let error_strings: Vec = errors.iter().map(std::string::ToString::to_string).collect(); diff --git a/src/model/api_proxy.rs b/src/model/api_proxy.rs index 1d8f75bb0..2c5e39bb9 100644 --- a/src/model/api_proxy.rs +++ b/src/model/api_proxy.rs @@ -174,6 +174,7 @@ impl FromStr for ProxyUserStatus { fn from_str(s: &str) -> Result { match s { + Self::ACTIVE => Ok(Self::Active), Self::EXPIRED => Ok(Self::Expired), Self::BANNED => Ok(Self::Banned), Self::TRIAL => Ok(Self::Trial), diff --git a/src/processing/processor/playlist.rs b/src/processing/processor/playlist.rs index 307fdc307..241999244 100644 --- a/src/processing/processor/playlist.rs +++ b/src/processing/processor/playlist.rs @@ -261,7 +261,7 @@ async fn process_source(client: Arc, cfg: Arc, source_i let start_time = Instant::now(); let (mut playlistgroups, mut error_list) = match input.input_type { InputType::M3u => m3u::get_m3u_playlist(Arc::clone(&client), &cfg, input, &cfg.working_dir).await, - InputType::Xtream => xtream::get_xtream_playlist(Arc::clone(&client), input, &cfg.working_dir).await, + InputType::Xtream => xtream::get_xtream_playlist(&cfg, Arc::clone(&client), input, &cfg.working_dir).await, InputType::M3uBatch | InputType::XtreamBatch => (vec![], vec![]) }; let (tvguide, mut tvguide_errors) = if error_list.is_empty() { diff --git a/src/utils/json_utils.rs b/src/utils/json_utils.rs index ca57cc9a6..23df83108 100644 --- a/src/utils/json_utils.rs +++ b/src/utils/json_utils.rs @@ -174,6 +174,14 @@ pub fn get_u64_from_serde_value(value: &Value) -> Option { } } +pub fn get_i64_from_serde_value(value: &Value) -> Option { + match value { + Value::Number(num_val) => num_val.as_i64(), + Value::String(str_val) => str_val.parse::().ok(), + _ => None, + } +} + pub fn get_u32_from_serde_value(value: &Value) -> Option { get_u64_from_serde_value(value).and_then(|val| u32::try_from(val).ok()) } diff --git a/src/utils/network/xtream.rs b/src/utils/network/xtream.rs index a6444d81f..2b259a3cd 100644 --- a/src/utils/network/xtream.rs +++ b/src/utils/network/xtream.rs @@ -1,18 +1,21 @@ -use std::sync::Arc; -use crate::tuliprox_error::{str_to_io_error, TuliproxError}; -use crate::model::{Config, ConfigInput, ConfigTarget}; +use crate::model::ProxyUserCredentials; +use crate::model::{Config, ConfigInput, ConfigTarget, ProxyUserStatus}; use crate::model::{PlaylistEntry, PlaylistGroup, XtreamCluster, XtreamPlaylistItem}; use crate::processing::parser::xtream; -use crate::repository::xtream_repository::{rewrite_xtream_series_info_content, rewrite_xtream_vod_info_content, xtream_get_input_info}; use crate::repository::xtream_repository; +use crate::repository::xtream_repository::{rewrite_xtream_series_info_content, rewrite_xtream_vod_info_content, xtream_get_input_info}; +use crate::tuliprox_error::{str_to_io_error, TuliproxError}; +use crate::utils; +use crate::utils::{get_string_from_serde_value, request}; +use crate::utils::request::extract_extension_from_url; +use chrono::{DateTime}; use log::{info, warn}; use std::cmp::Ordering; use std::io::Error; -use crate::model::{ProxyUserCredentials}; -use crate::utils::get_string_from_serde_value; -use crate::utils::request; -use crate::utils::request::extract_extension_from_url; - +use std::str::FromStr; +use std::sync::Arc; +use std::time::{SystemTime, UNIX_EPOCH}; +use crate::messaging::{send_message, MsgKind}; #[inline] pub fn get_xtream_stream_url_base(url: &str, username: &str, password: &str) -> String { @@ -129,13 +132,71 @@ const ACTIONS: [(XtreamCluster, &str, &str); 3] = [ (XtreamCluster::Video, crate::model::XC_ACTION_GET_VOD_CATEGORIES, crate::model::XC_ACTION_GET_VOD_STREAMS), (XtreamCluster::Series, crate::model::XC_ACTION_GET_SERIES_CATEGORIES, crate::model::XC_ACTION_GET_SERIES)]; -pub async fn get_xtream_playlist(client: Arc, input: &ConfigInput, working_dir: &str) -> (Vec, Vec) { +async fn xtream_login(cfg: &Config, client: &Arc, input: &ConfigInput, username: &str, base_url: &str) -> Result<(), TuliproxError> { + let content = match request::get_input_json_content(Arc::clone(client), input, base_url, None).await { + Ok(content) => content, + Err(_) => { + match request::get_input_json_content(Arc::clone(client), input, &format!("{base_url}&action=get_account_info"), None).await { + Ok(content) => content, + Err(err) => { + warn!("Failed to login xtream account {username} {err}"); + return Err(err); + } + } + } + }; + match content.get("user_info") { + None => {} + Some(value) => { + if let Some(status_value) = value.get("status") { + if let Some(status) = utils::get_string_from_serde_value(status_value) { + if let Ok(cur_status) = ProxyUserStatus::from_str(&status) { + if matches!(cur_status, ProxyUserStatus::Active | ProxyUserStatus::Trial) { + warn!("User status for user {username} is {cur_status:?}"); + send_message(client, &MsgKind::Info, cfg.messaging.as_ref(), &format!("User status for user {username} is {cur_status:?}")); + } + } + } + } + if let Some(status_value) = value.get("exp_date") { + if let Some(expiration_timestamp) = utils::get_i64_from_serde_value(status_value) { + if expiration_timestamp > 0 { + #[allow(clippy::cast_sign_loss)] + let expiration_ts = expiration_timestamp as u64; + if let Ok(now) = SystemTime::now().duration_since(UNIX_EPOCH) { + let now_secs = now.as_secs(); + if expiration_ts > now_secs { + let time_left = expiration_ts - now_secs; + if time_left > 3 * 24 * 60 * 60 { + let datetime = DateTime::from_timestamp(expiration_timestamp, 0).unwrap(); + let formatted = datetime.format("%Y-%m-%d %H:%M:%S").to_string(); + warn!("User account for user {username} expires {formatted}"); + send_message(client, &MsgKind::Info, cfg.messaging.as_ref(), &format!("User account for user {username} expires {formatted}")); + } + } else { + warn!("User account for user {username} is expired"); + send_message(client, &MsgKind::Info, cfg.messaging.as_ref(), &format!("User account for user {username} is expired")); + } + } + } + } + } + } + } + Ok(()) +} + +pub async fn get_xtream_playlist(cfg: &Config, client: Arc, input: &ConfigInput, working_dir: &str) -> (Vec, Vec) { let username = input.username.as_ref().map_or("", |v| v); let password = input.password.as_ref().map_or("", |v| v); let base_url = get_xtream_stream_url_base(&input.url, username, password); + if let Err(err) = xtream_login(cfg, &client, input, username, &base_url).await { + return (Vec::with_capacity(0), vec![err]); + } + if let Err(_err) = request::get_input_json_content(Arc::clone(&client), input, base_url.as_str(), None).await { if let Err(err) = request::get_input_json_content(Arc::clone(&client), input, &format!("{base_url}&action=get_account_info"), None).await { warn!("Failed to login xtream account {username} {err}"); @@ -188,10 +249,10 @@ pub async fn get_xtream_playlist(client: Arc, input: &ConfigInp pub fn create_vod_info_from_item(target: &ConfigTarget, user: &ProxyUserCredentials, pli: &XtreamPlaylistItem, last_updated: i64) -> String { let category_id = pli.category_id; - let stream_id = if user.proxy.is_redirect(pli.item_type) || target.is_force_redirect(pli.item_type) { pli.provider_id } else { pli.virtual_id }; + let stream_id = if user.proxy.is_redirect(pli.item_type) || target.is_force_redirect(pli.item_type) { pli.provider_id } else { pli.virtual_id }; let name = &pli.name; let extension = pli.get_additional_property("container_extension") - .map_or_else(|| extract_extension_from_url(&pli.url).map_or_else (String::new, std::string::ToString::to_string), + .map_or_else(|| extract_extension_from_url(&pli.url).map_or_else(String::new, std::string::ToString::to_string), |v| get_string_from_serde_value(&v).map_or_else(String::new, |v| v)); let added = last_updated / 1000; format!(r#"{{