Merge pull request #455 from euzu/feature/refactor_provider_connection_handling

- Short EPG is now served from local disk, if available
- WebUI Api-User Category selection implemented
- Stream Table Copy-To-Clipboard functions added
- Refactored provider connection handling to avoid possible race conditions
- Added `exp_date` field to inputs, aliases, and CSV batch files; accepts date in "YYYY-MM-DD HH:MM:SS" format or Unix timestamp (seconds since epoch).
This commit is contained in:
euzu
2025-11-26 14:35:24 +01:00
committed by GitHub
47 changed files with 1347 additions and 409 deletions
+5 -1
View File
@@ -23,7 +23,11 @@
- Added `order: none` support for group/channel sorting so mappings can opt out of any reordering and keep the source order.
- Session tracking now matches repeated HLS segment connections by session token so a single user keeps one active connection count even when new TCP sockets are opened.
- EPG icon urls are now rewritten on reverse proxy mode.
- Xtream Codes Batch provider accounts are now checked for expiration.
- Short EPG is now served from local disk, if available
- WebUI Api-User Category selection implemented
- Stream Table Copy-To-Clipboard functions added
- Refactored provider connection handling to avoid possible race conditions
- Added `exp_date` field to inputs, aliases, and CSV batch files; accepts date in "YYYY-MM-DD HH:MM:SS" format or Unix timestamp (seconds since epoch).
# 3.2.0 (2025-11-14)
- Added `name` attribute to Staged Input.
Generated
+3 -3
View File
@@ -1096,7 +1096,7 @@ dependencies = [
[[package]]
name = "frontend"
version = "3.2.10"
version = "3.2.11"
dependencies = [
"anyhow",
"base64",
@@ -3765,7 +3765,7 @@ dependencies = [
[[package]]
name = "shared"
version = "3.2.10"
version = "3.2.11"
dependencies = [
"base64",
"bitflags 2.10.0",
@@ -4314,7 +4314,7 @@ checksum = "e421abadd41a4225275504ea4d6566923418b7f05506fbc9c0fe86ba7396114b"
[[package]]
name = "tuliprox"
version = "3.2.10"
version = "3.2.11"
dependencies = [
"arc-swap",
"async-compression",
+10 -7
View File
@@ -665,9 +665,6 @@ The template can now be used for sequence
- '(?i)\bSD\b'
```
### 2.2. `sources`
`sources` is a sequence of source definitions, which have two top level entries:
-`inputs`
@@ -687,7 +684,8 @@ Each input has the following attributes:
- `headers` is optional
- `method` can be `GET` or `POST`
- `username` only mandatory for type `xtream`
- `pasword`only mandatory for type `xtream`
- `password` only mandatory for type `xtream`
- `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)
- `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.
@@ -834,9 +832,9 @@ There are 2 batch input types `xtream_batch` and `m3u_batch`.
```
```csv
#name;username;password;url;max_connections;priority
my_provider_1;user1;password1;http://my_provider_1.com:80;1;0
my_provider_2;user2;password2;http://my_provider_2.com:8080;1;0
#name;username;password;url;max_connections;priority;exp_date
my_provider_1;user1;password1;http://my_provider_1.com:80;1;0;2028-11-23 12:34:23
my_provider_2;user2;password2;http://my_provider_2.com:8080;1;0;2028-11-23 12:34:23
```
##### `M3uBatch`
@@ -864,6 +862,11 @@ A `priority` of `0` is higher than `1`
Higher numbers mean **lower priority**
This means tasks or items with smaller (even negative) values will be handled before those with larger values.
The `exp_date` field is a date as:
- "YYYY-MM-DD HH:MM:SS" format like `2028-11-30 12:34:12`
- or Unix timestamp (seconds since epoch)
### 2.2.2 `targets`
Has the following top level entries:
- `enabled` _optional_ default is `true`, if you disable the processing is skipped
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "tuliprox"
version = "3.2.10"
version = "3.2.11"
edition = "2021"
# See more keys and their definitions at https://doc.rust-lang.org/cargo/reference/manifest.html
+2 -2
View File
@@ -360,7 +360,7 @@ async fn resolve_streaming_strategy(
if let Some(allocation) = provider_connection_handle.as_ref().map(|ph| &ph.allocation) {
match allocation {
ProviderAllocation::Exhausted => {
debug!("Input {} is exhausted. No connections allowed.", input.name);
debug!("Provider {} is exhausted. No connections allowed.", input.name);
let stream = create_provider_connections_exhausted_stream(&app_state.app_config, &[]);
ProviderStreamState::Custom(stream)
}
@@ -391,7 +391,7 @@ async fn resolve_streaming_strategy(
}
}
} else {
debug!("Input {} is exhausted. No connections allowed.", input.name);
debug!("Provider {} is exhausted. No connections allowed.", input.name);
let stream = create_provider_connections_exhausted_stream(&app_state.app_config, &[]);
ProviderStreamState::Custom(stream)
};
+5 -3
View File
@@ -166,10 +166,11 @@ pub(in crate::api) async fn handle_hls_stream_request(
async fn get_stream_channel(app_state: &Arc<AppState>, target: &Arc<ConfigTarget>, virtual_id: u32) -> Option<StreamChannel> {
if target.has_output(TargetType::Xtream) {
if let Ok((pli, _)) = xtream_repository::xtream_get_item_for_stream_id(virtual_id, app_state, target, None).await {
return Some(pli.to_stream_channel());
return Some(pli.to_stream_channel(target.id));
}
}
m3u_get_item_for_stream_id(virtual_id, app_state, target).await.ok().map(|pli| pli.to_stream_channel())
let target_id = target.id;
m3u_get_item_for_stream_id(virtual_id, app_state, target).await.ok().map(|pli| pli.to_stream_channel(target_id))
}
async fn resolve_stream_channel(
@@ -181,6 +182,7 @@ async fn resolve_stream_channel(
let mut channel = match get_stream_channel(app_state, target, virtual_id).await {
Some(channel) => channel,
None => StreamChannel {
target_id: target.id,
virtual_id,
provider_id: 0,
item_type: PlaylistItemType::LiveHls,
@@ -225,7 +227,7 @@ async fn hls_api_stream(
app_state.app_config.get_input_by_id(params.input_id),
true,
format!(
"Cant find input for target {target_name}, stream_id {virtual_id}, hls"
"Cant find input {} for target {target_name}, stream_id {virtual_id}, hls", params.input_id
)
);
+3 -3
View File
@@ -119,7 +119,7 @@ async fn m3u_api_stream(
.app_config
.get_input_by_name(pli.input_name.as_str()),
true,
format!("Cant find input for target {target_name}, stream_id {virtual_id}")
format!("Cant find input {} for target {target_name}, stream_id {virtual_id}", pli.input_name)
);
let cluster = XtreamCluster::try_from(pli.item_type).unwrap_or(XtreamCluster::Live);
@@ -154,7 +154,7 @@ async fn m3u_api_stream(
fingerprint,
app_state,
session,
pli.to_stream_channel(),
pli.to_stream_channel(target.id),
req_headers,
&input,
&user,
@@ -224,7 +224,7 @@ async fn m3u_api_stream(
fingerprint,
app_state,
&session_key,
pli.to_stream_channel(),
pli.to_stream_channel(target.id),
session_url,
req_headers,
&input,
+84 -27
View File
@@ -145,11 +145,12 @@ fn parse_timeshift(time_shift: Option<&String>) -> Option<i32> {
})
}
async fn serve_epg(
pub async fn serve_epg(
app_state: &Arc<AppState>,
epg_path: &Path,
user: &ProxyUserCredentials,
target: &Arc<ConfigTarget>,
filter: Option<String>,
) -> axum::response::Response {
if let Ok(exists) = tokio::fs::try_exists(epg_path).await {
if exists {
@@ -165,12 +166,12 @@ async fn serve_epg(
// Use 0 for timeshift if None
let timeshift = parse_timeshift(user.epg_timeshift.as_ref()).unwrap_or(0);
return if timeshift != 0 || rewrite_urls {
return if timeshift != 0 || rewrite_urls || filter.is_some() {
let server_info = app_state.app_config.get_user_server_info(user);
let base_url = format!("{}/{}/{}/{}/", server_info.get_base_url(),
storage_const::EPG_RESOURCE_PATH, &user.username, &user.password);
// Apply timeshift and/or rewrite URLs
serve_epg_with_rewrites(epg_path, timeshift, rewrite_urls, &encrypt_secret, &base_url).await
// Apply timeshift and/or rewrite URLs and/or filter
serve_epg_with_rewrites(epg_path, timeshift, rewrite_urls, &encrypt_secret, &base_url, filter).await
} else {
// Neither timeshift nor rewrite needed, serve original file
serve_file(epg_path, mime::TEXT_XML).await.into_response()
@@ -187,6 +188,7 @@ async fn serve_epg_with_rewrites(
rewrite_urls: bool,
secret: &[u8; 16],
base_url: &str,
filter: Option<String>,
) -> axum::response::Response {
match tokio::fs::try_exists(epg_path).await {
Ok(exists) => {
@@ -231,10 +233,75 @@ async fn serve_epg_with_rewrites(
let mut buf = Vec::with_capacity(4096);
let duration = Duration::minutes(i64::from(offset_minutes));
let mut skip_depth = None;
loop {
match xml_reader.read_event_into_async(&mut buf).await {
Ok(Event::Start(ref e)) if offset_minutes != 0 && e.name().as_ref() == b"programme" => {
buf.clear();
let event = match xml_reader.read_event_into_async(&mut buf).await {
Ok(e) => e,
Err(e) => {
error!("Error reading epg XML event: {e}");
break;
}
};
if let Some(flt) = &filter {
// Filter
match &event {
Event::Start(e) => {
if skip_depth.is_none() {
let should_skip = match e.name().as_ref() {
b"channel" => {
e.attributes()
.filter_map(Result::ok)
.find(|a| a.key.as_ref() == b"id")
.and_then(|a| a.unescape_value().ok())
.is_some_and(|v| !flt.eq(v.as_ref()))
}
b"programme" => {
e.attributes()
.filter_map(Result::ok)
.find(|a| a.key.as_ref() == b"channel")
.and_then(|a| a.unescape_value().ok())
.is_some_and(|v| !flt.eq(v.as_ref()))
}
_ => false,
};
if should_skip {
skip_depth = Some(1);
continue;
}
} else {
skip_depth = skip_depth.map(|d| d + 1);
continue;
}
}
Event::End(_) => {
if let Some(depth) = skip_depth {
if depth == 1 {
skip_depth = None;
} else {
skip_depth = Some(depth - 1);
}
continue;
}
}
Event::Empty(_) => {
if skip_depth.is_some() {
continue;
}
}
_ => {}
}
if skip_depth.is_some() {
continue;
}
}
match &event {
Event::Start(ref e) if offset_minutes != 0 && e.name().as_ref() == b"programme" => {
// Modify the attributes
let mut elem = BytesStart::new(EPG_TAG_PROGRAMME);
for attr in e.attributes() {
@@ -261,18 +328,18 @@ async fn serve_epg_with_rewrites(
elem.push_attribute(attr);
}
Err(e) => {
error!("Error parsing attribute: {e}");
error!("Error parsing epg attribute: {e}");
}
}
}
// Write the modified start event
if let Err(e) = xml_writer.write_event_async(Event::Start(elem)).await {
error!("Failed to write Start event: {e}");
error!("Failed to write epg Start event: {e}");
break;
}
}
Ok(ref event @ (Event::Empty(ref e) | Event::Start(ref e))) if rewrite_urls && e.name().as_ref() == b"icon" => {
ref event @ (Event::Empty(ref e) | Event::Start(ref e)) if rewrite_urls && e.name().as_ref() == b"icon" => {
// Modify the attributes
let mut elem = BytesStart::new(EPG_TAG_ICON);
for attr in e.attributes() {
@@ -298,16 +365,11 @@ async fn serve_epg_with_rewrites(
elem.push_attribute(attr);
}
Err(e) => {
error!("Error parsing attribute: {e}");
error!("Error parsing epg attribute: {e}");
}
}
}
// Write the modified icon event
// if let Err(e) = xml_writer.write_event_async(Event::Start(elem)).await {
// error!("Failed to write Start event: {e}");
// break;
// }
let out_event = match event {
Event::Empty(_) => Some(Event::Empty(elem)),
Event::Start(_) => Some(Event::Start(elem)),
@@ -315,28 +377,23 @@ async fn serve_epg_with_rewrites(
};
if let Some(out) = out_event {
if let Err(e) = xml_writer.write_event_async(out).await {
error!("Failed to write icon event: {e}");
error!("Failed to write epg icon event: {e}");
break;
}
}
}
Ok(Event::Decl(_) | Event::DocType(_)) => {},
Ok(Event::Eof) => break, // End of file
Ok(event) => {
Event::Decl(_) | Event::DocType(_) => {},
Event::Eof => break, // End of file
_ => {
// Write any other event as is
if let Err(e) = xml_writer.write_event_async(event).await {
error!("Failed to write event: {e}");
error!("Failed to epg write event: {e}");
break;
}
}
Err(e) => {
error!("Error: {e}");
break;
}
}
buf.clear();
}
buf.clear();
let mut encoder = xml_writer.into_inner();
if let Err(e) = encoder.shutdown().await {
error!("Failed to shutdown epg gzip encoder: {e}");
@@ -386,7 +443,7 @@ async fn xmltv_api(
return get_empty_epg_response();
};
serve_epg(&app_state, &epg_path, &user, &target).await
serve_epg(&app_state, &epg_path, &user, &target, None).await
}
#[axum::debug_handler]
+20 -10
View File
@@ -9,7 +9,7 @@ use crate::api::api_utils::{
};
use crate::api::api_utils::{redirect, try_result_not_found, try_option_bad_request, try_result_bad_request};
use crate::api::endpoints::hls_api::handle_hls_stream_request;
use crate::api::endpoints::xmltv_api::get_empty_epg_response;
use crate::api::endpoints::xmltv_api::{get_empty_epg_response, get_epg_path_for_target, serve_epg};
use crate::api::model::AppState;
use crate::api::model::UserApiRequest;
use crate::api::model::XtreamAuthorizationResponse;
@@ -256,7 +256,7 @@ async fn xtream_player_api_stream(
let input = try_option_bad_request!(
app_state.app_config.get_input_by_name(pli.input_name.as_str()),
true,
format!( "Cant find input for target {target_name}, context {}, stream_id {virtual_id}", stream_req.context)
format!( "Cant find input {} for target {target_name}, context {}, stream_id {virtual_id}", pli.input_name, stream_req.context)
);
let (cluster, item_type) = if stream_req.context == ApiStreamContext::Timeshift {
@@ -291,7 +291,7 @@ async fn xtream_player_api_stream(
.into_response();
}
let stream_channel = create_stream_channel_with_type(&pli, item_type);
let stream_channel = create_stream_channel_with_type(target.id, &pli, item_type);
if session.virtual_id == virtual_id && is_seek_request(cluster, req_headers).await {
// partial request means we are in reverse proxy mode, seek happened
@@ -380,7 +380,7 @@ async fn xtream_player_api_stream(
.into_response();
}
let stream_channel = create_stream_channel_with_type(&pli, item_type);
let stream_channel = create_stream_channel_with_type(target.id, &pli, item_type);
stream_response(
fingerprint,
@@ -432,8 +432,8 @@ async fn xtream_player_api_stream_with_token(
.get_input_by_name(pli.input_name.as_str()),
true,
format!(
"Cant find input for target {target_name}, context {}, stream_id {}",
stream_req.context, pli.virtual_id
"Cant find input {} for target {target_name}, context {}, stream_id {}",
pli.input_name, stream_req.context, pli.virtual_id
)
);
@@ -516,7 +516,7 @@ async fn xtream_player_api_stream_with_token(
fingerprint,
app_state,
session_key.as_str(),
pli.to_stream_channel(),
pli.to_stream_channel(target.id),
&stream_url,
req_headers,
&input,
@@ -946,7 +946,7 @@ async fn xtream_player_api_timeshift_query_stream(
async fn xtream_get_stream_info_response(
app_state: &Arc<AppState>,
user: &ProxyUserCredentials,
target: &ConfigTarget,
target: &Arc<ConfigTarget>,
stream_id: &str,
cluster: XtreamCluster,
) -> impl IntoResponse + Send {
@@ -1029,7 +1029,7 @@ async fn xtream_get_stream_info_response(
async fn xtream_get_short_epg(
app_state: &Arc<AppState>,
user: &ProxyUserCredentials,
target: &ConfigTarget,
target: &Arc<ConfigTarget>,
stream_id: &str,
limit: &str,
) -> impl IntoResponse + Send {
@@ -1046,6 +1046,16 @@ async fn xtream_get_short_epg(
target,
None,
).await {
let config = &app_state.app_config.config.load();
if let Some(epg_path) = get_epg_path_for_target(config, target) {
if let Ok(exists) = tokio::fs::try_exists(&epg_path).await {
if exists {
return serve_epg(app_state, &epg_path, user, target, pli.epg_channel_id.clone()).await
}
}
}
if pli.provider_id > 0 {
let input_name = &pli.input_name;
if let Some(input) = app_state.app_config.get_input_by_name(input_name.as_str()) {
@@ -1196,7 +1206,7 @@ async fn xtream_player_api_handle_content_action(
async fn xtream_get_catchup_response(
app_state: &Arc<AppState>,
target: &ConfigTarget,
target: &Arc<ConfigTarget>,
stream_id: &str,
start: &str,
end: &str,
@@ -2,41 +2,42 @@ use crate::api::model::provider_lineup_manager::{ProviderAllocation, ProviderLin
use crate::api::model::{EventManager, ProviderConfig};
use crate::model::{AppConfig, ConfigInput};
use log::{debug, error};
use crate::utils::{trace_if_enabled};
use shared::utils::{default_grace_period_millis, default_grace_period_timeout_secs};
use std::collections::{HashMap, HashSet};
use std::net::SocketAddr;
use std::sync::Arc;
use tokio::sync::RwLock;
pub type ProviderConnectionId = SocketAddr;
pub type ClientConnectionId = SocketAddr;
#[derive(Debug, Clone)]
pub struct ProviderHandle {
pub id: ProviderConnectionId,
pub client_id: ClientConnectionId,
pub allocation: ProviderAllocation,
}
impl ProviderHandle {
pub fn new(id: ProviderConnectionId, allocation: ProviderAllocation) -> Self {
Self { id, allocation }
pub fn new(client_id: ClientConnectionId, allocation: ProviderAllocation) -> Self {
Self { client_id, allocation }
}
}
#[derive(Debug, Clone)]
struct SharedAllocation {
allocation: ProviderAllocation,
connections: HashSet<ProviderConnectionId>,
connections: HashSet<ClientConnectionId>,
}
#[derive(Debug, Clone, Default)]
struct SharedConnections {
by_key: HashMap<String, SharedAllocation>,
key_by_addr: HashMap<SocketAddr, String>,
key_by_addr: HashMap<ClientConnectionId, String>,
}
#[derive(Debug, Clone, Default)]
struct Connections {
single: HashMap<ProviderConnectionId, ProviderAllocation>,
single: HashMap<ClientConnectionId, ProviderAllocation>,
shared: SharedConnections,
}
@@ -74,9 +75,6 @@ impl ActiveProviderManager {
}
async fn acquire_connection_inner(&self, provider_or_input_name: &str, addr: &SocketAddr, force: bool) -> Option<ProviderHandle> {
// Lock connections
let mut connections = self.connections.write().await;
// Call the specific acquisition function
let allocation = if force {
self.providers.force_exact_acquire_connection(provider_or_input_name).await
@@ -88,12 +86,13 @@ impl ActiveProviderManager {
ProviderAllocation::Exhausted => {}
ProviderAllocation::Available(_) | ProviderAllocation::GracePeriod(_) => {
let provider_name = allocation.get_provider_name().unwrap_or_default();
let mut connections = self.connections.write().await;
if let Some(old) = connections.single.insert(*addr, allocation.clone()) {
crate::utils::trace_if_enabled!(
trace_if_enabled!(
"register_connection: address {addr} already had a allocation for provider {:?} — forcing release on the old allocation",
old.get_provider_name().unwrap_or_default()
);
old.get_provider_name().unwrap_or_default());
drop(connections);
old.release().await;
}
@@ -128,37 +127,64 @@ impl ActiveProviderManager {
}
pub async fn release_connection(&self, addr: &SocketAddr) {
let mut connections = self.connections.write().await;
// Single connection
let single_allocation = {
let mut connections = self.connections.write().await;
connections.single.remove(addr)
};
// try to release the single connection (not shared)
let handle = connections.single.remove(addr);
if let Some(allocation) = handle {
debug!("Released provider connection {:?} for {addr}", allocation.get_provider_name().unwrap_or_default());
if let Some(allocation) = single_allocation {
debug!(
"Released provider connection {:?} for {addr}",
allocation.get_provider_name().unwrap_or_default()
);
allocation.release().await;
return;
}
let key = match connections.shared.key_by_addr.get(addr) {
Some(k) => k.clone(),
None => return,
// Shared connection
let shared_allocation = {
let mut connections = self.connections.write().await;
let key = match connections.shared.key_by_addr.get(addr) {
Some(k) => k.clone(),
None => return, // no shared connection
};
// Clone the SharedAllocation to avoid double mutable borrow
let mut shared = match connections.shared.by_key.get(&key) {
Some(s) => s.clone(),
None => return,
};
// Remove this address from the shared connection set
shared.connections.remove(addr);
// Always remove stale key-by-addr entry
connections.shared.key_by_addr.remove(addr);
if shared.connections.is_empty() {
// If this was the last user of the shared allocation:
connections.shared.by_key.remove(&key);
Some(shared.allocation)
} else {
// Update the entry back with the remaining connections
connections.shared.by_key.insert(key, shared);
None
}
};
let mut released = false;
if let Some(connections) = connections.shared.by_key.get_mut(&key) {
if connections.connections.remove(addr) && connections.connections.is_empty() {
connections.allocation.release().await;
released = true;
}
}
if released {
connections.shared.key_by_addr.remove(addr);
connections.shared.by_key.remove(&key);
// release allocation
if let Some(allocation) = shared_allocation {
allocation.release().await;
debug!(
"Released last shared connection for provider {}, releasing allocation {addr}",
allocation.get_provider_name().unwrap_or_default()
);
}
}
pub async fn release_handle(&self, handle: &ProviderHandle) {
self.release_connection(&handle.id).await;
self.release_connection(&handle.client_id).await;
}
pub async fn make_shared_connection(&self, addr: &SocketAddr, key: &str) {
@@ -178,7 +204,7 @@ impl ActiveProviderManager {
shared_allocation.connections.insert(*addr);
connections.shared.key_by_addr.insert(*addr, key.to_string());
} else {
error!("Failed to add shared connection for {addr}: url: {key:?} not found");
error!("Failed to add shared connection for {addr}: url: {key:?} not found");
}
}
+1 -1
View File
@@ -58,7 +58,7 @@ impl ConnectionManager {
pub async fn release_provider_handle(&self, provider_handle: Option<ProviderHandle>) {
if let Some(handle) = provider_handle {
self.release_provider_connection(&handle.id).await;
self.release_provider_connection(&handle.client_id).await;
}
}
+26 -6
View File
@@ -1,5 +1,5 @@
use std::fmt;
use crate::model::{ConfigInput, ConfigInputAlias, InputUserInfo};
use crate::model::{is_input_expired, ConfigInput, ConfigInputAlias, InputUserInfo};
use jsonwebtoken::get_current_timestamp;
use log::{debug};
use std::ops::Deref;
@@ -42,6 +42,7 @@ pub struct ProviderConfig {
pub input_type: InputType,
max_connections: usize,
priority: i16,
exp_date: Option<i64>,
connection: RwLock<ProviderConfigConnection>,
on_connection_change: ProviderConnectionChangeCallback,
}
@@ -57,7 +58,8 @@ impl fmt::Display for ProviderConfig {
write!(f, ", priority: {}", self.priority)?;
write_if_some!(f, self,
", username: " => username,
", password: " => password
", password: " => password,
", exp_date: " => exp_date
);
write!(f, "}}")?;
Ok(())
@@ -80,6 +82,7 @@ impl PartialEq for ProviderConfig {
&& self.input_type == other.input_type
&& self.max_connections == other.max_connections
&& self.priority == other.priority
&& self.exp_date == other.exp_date
// Note: self.connection is skipped
}
}
@@ -90,7 +93,7 @@ macro_rules! modify_connections {
$self.notify_connection_change($guard.current_connections);
}};
($self:ident, $guard:ident, -1) => {{
$guard.current_connections -= 1;
$guard.current_connections = $guard.current_connections.saturating_sub(1);
$self.notify_connection_change($guard.current_connections);
}};
}
@@ -109,6 +112,7 @@ impl ProviderConfig {
input_type: cfg.input_type,
max_connections: cfg.max_connections as usize,
priority: cfg.priority,
exp_date: cfg.exp_date,
connection: RwLock::new(get_connection.and_then(|f| f(cfg.name.as_str())).map_or_else(Default::default, Clone::clone)),
on_connection_change
}
@@ -127,6 +131,7 @@ impl ProviderConfig {
input_type: cfg.input_type,
max_connections: alias.max_connections as usize,
priority: alias.priority,
exp_date: alias.exp_date,
connection: RwLock::new(get_connection.and_then(|f| f(alias.name.as_str())).map_or_else(Default::default, Clone::clone)),
on_connection_change,
}
@@ -178,12 +183,20 @@ impl ProviderConfig {
// !self.is_exhausted()
// }
async fn force_allocate(&self) {
async fn force_allocate(&self) -> bool {
if is_input_expired(self.exp_date) {
return false;
}
let mut guard = self.connection.write().await;
modify_connections!(self, guard, +1);
true
}
async fn try_allocate(&self, grace: bool, grace_period_timeout_secs: u64) -> ProviderConfigAllocation {
if is_input_expired(self.exp_date) {
return ProviderConfigAllocation::Exhausted;
}
let mut guard = self.connection.write().await;
if self.max_connections == 0 {
modify_connections!(self, guard, +1);
@@ -220,6 +233,10 @@ impl ProviderConfig {
// is intended to use with redirects, to cycle through provider
// do not increment and connection counter!
async fn get_next(&self, grace: bool, grace_period_timeout_secs: u64) -> bool {
if is_input_expired(self.exp_date) {
return false;
}
if self.max_connections == 0 {
return true;
}
@@ -288,8 +305,11 @@ impl ProviderConfigWrapper {
}
pub async fn force_allocate(&self) -> ProviderAllocation {
self.inner.force_allocate().await;
ProviderAllocation::new_available(Arc::clone(&self.inner))
if self.inner.force_allocate().await {
ProviderAllocation::new_available(Arc::clone(&self.inner))
} else {
ProviderAllocation::Exhausted
}
}
pub async fn try_allocate(&self, grace: bool, grace_period_timeout_secs: u64) -> ProviderAllocation {
@@ -555,6 +555,7 @@ impl ProviderLineupManager {
|| a.username != b.username
|| a.password != b.password
|| a.url != b.url
|| a.exp_date != b.exp_date
{
return true;
}
@@ -577,6 +578,7 @@ impl ProviderLineupManager {
|| a_alias.username != b_alias.username
|| a_alias.password != b_alias.password
|| a_alias.url != b_alias.url
|| a_alias.exp_date != b_alias.exp_date
{
return true;
}
@@ -1,4 +1,4 @@
use crate::api::model::{AppState, CustomVideoStreamType, ProviderHandle, StreamDetails};
use crate::api::model::{AppState, ConnectionManager, CustomVideoStreamType, ProviderHandle, StreamDetails};
use crate::api::model::BoxedProviderStream;
use crate::api::model::StreamError;
use crate::api::model::TimedClientStream;
@@ -30,6 +30,7 @@ pub(in crate::api) struct ActiveClientStream {
provider_handle: Option<ProviderHandle>,
custom_video: (Option<TransportStreamBuffer>, Option<TransportStreamBuffer>),
waker: Option<Arc<AtomicWaker>>,
connection_manager: Arc<ConnectionManager>,
}
impl ActiveClientStream {
@@ -95,6 +96,7 @@ impl ActiveClientStream {
send_custom_stream_flag: grace_stop_flag,
custom_video,
waker,
connection_manager: Arc::clone(&app_state.connection_manager)
}
}
@@ -218,3 +220,13 @@ impl Stream for ActiveClientStream {
Poll::Ready(None)
}
}
impl Drop for ActiveClientStream {
fn drop(&mut self) {
let mgr = Arc::clone(&self.connection_manager);
let hndl = self.provider_handle.take();
tokio::spawn(async move {
mgr.release_provider_handle(hndl).await;
});
}
}
@@ -354,11 +354,8 @@ async fn handle_channel_unavailable_stream(app_state: &Arc<AppState>,
app_state.connection_manager.release_provider_connection(&stream_options.addr).await;
if let (Some(boxed_provider_stream), response_info) =
create_channel_unavailable_stream(
&app_state.app_config,
&get_response_headers(stream_options.get_headers()),
StatusCode::SERVICE_UNAVAILABLE,
)
create_channel_unavailable_stream(&app_state.app_config,&get_response_headers(stream_options.get_headers()),
StatusCode::SERVICE_UNAVAILABLE)
{
Ok(Some((boxed_provider_stream, response_info)))
} else {
@@ -383,37 +380,23 @@ async fn get_provider_stream(
}
Ok(None) => {
if connect_err > ERR_MAX_RETRY_COUNT {
warn!(
"The stream could be unavailable. {}",
sanitize_sensitive_info(stream_options.get_url().as_str())
);
warn!("The stream could be unavailable. {}", sanitize_sensitive_info(stream_options.get_url().as_str()));
break;
}
}
Err(status) => {
debug!("Provider stream response error status response : {status}");
if status == StatusCode::FORBIDDEN
|| status == StatusCode::SERVICE_UNAVAILABLE
|| status == StatusCode::UNAUTHORIZED
{
warn!(
"The stream could be unavailable. ({status}) {}",
sanitize_sensitive_info(stream_options.get_url().as_str())
);
stream_options.cancel_reconnect();
return Err(status);
if matches!(status, StatusCode::FORBIDDEN | StatusCode::SERVICE_UNAVAILABLE | StatusCode::UNAUTHORIZED) {
warn!("The stream could be unavailable. ({status}) {}",sanitize_sensitive_info(stream_options.get_url().as_str()));
break;
}
if connect_err > ERR_MAX_RETRY_COUNT {
warn!(
"The stream could be unavailable. ({status}) {}",
sanitize_sensitive_info(stream_options.get_url().as_str())
);
warn!("The stream could be unavailable. ({status}) {}",sanitize_sensitive_info(stream_options.get_url().as_str()));
break;
}
}
}
if !stream_options.should_continue() {
return Err(StatusCode::SERVICE_UNAVAILABLE);
}
if connect_err > ERR_MAX_RETRY_COUNT {
if !stream_options.should_continue() || connect_err > ERR_MAX_RETRY_COUNT {
break;
}
if start.elapsed().as_secs() > RETRY_SECONDS {
@@ -425,16 +408,11 @@ async fn get_provider_stream(
}
connect_err += 1;
tokio::time::sleep(Duration::from_millis(50)).await;
debug_if_enabled!(
"Reconnecting stream {}",
sanitize_sensitive_info(url.as_str())
);
debug_if_enabled!("Reconnecting stream {}", sanitize_sensitive_info(url.as_str()));
}
debug_if_enabled!(
"Stopped reconnecting stream {}",
sanitize_sensitive_info(url.as_str())
);
debug_if_enabled!("Stopped reconnecting stream {}", sanitize_sensitive_info(url.as_str()));
stream_options.cancel_reconnect();
app_state.connection_manager.release_provider_connection(&stream_options.addr).await;
Err(StatusCode::SERVICE_UNAVAILABLE)
}
@@ -493,7 +471,20 @@ pub async fn create_provider_stream(
if continue_streaming.is_active() {
match get_provider_stream(&app_state_clone, &client, &stream_opts).await {
Ok(Some((stream, _info))) => Some((stream, ())),
Ok(None) => None,
Ok(None) => {
app_state_clone.connection_manager.release_provider_connection(&stream_opts.addr).await;
continue_streaming.notify();
if let (Some(boxed_provider_stream), _response_info) =
create_channel_unavailable_stream(
&app_state_clone.app_config,
&get_response_headers(stream_opts.get_headers()),
StatusCode::SERVICE_UNAVAILABLE,
)
{
return Some((boxed_provider_stream, ()));
}
None
}
Err(status) => {
app_state_clone.connection_manager.release_provider_connection(&stream_opts.addr).await;
continue_streaming.notify();
@@ -510,6 +501,7 @@ pub async fn create_provider_stream(
}
}
} else {
app_state_clone.connection_manager.release_provider_connection(&stream_opts.addr).await;
None
}
}
@@ -20,6 +20,7 @@ use std::pin::Pin;
use std::task::{Context, Poll};
use tokio::sync::mpsc::Sender;
use tokio::sync::{mpsc, Mutex, RwLock};
use tokio::time::{sleep, Duration, Instant};
use tokio_stream::wrappers::ReceiverStream;
use tokio_util::sync::CancellationToken;
@@ -96,7 +97,7 @@ impl BurstBuffer {
}
pub fn push(&mut self, packet: Arc<Bytes>) {
while self.current_bytes > self.buffer_size {
while self.current_bytes + packet.len() > self.buffer_size {
if let Some(popped) = self.buffer.pop_front() {
self.current_bytes -= popped.len();
} else {
@@ -135,6 +136,7 @@ pub struct SharedStreamState {
broadcaster: tokio::sync::broadcast::Sender<Bytes>,
stop_token: CancellationToken,
burst_buffer: Arc<Mutex<BurstBuffer>>,
task_handles: Mutex<Vec<tokio::task::JoinHandle<()>>>,
}
impl SharedStreamState {
@@ -151,6 +153,7 @@ impl SharedStreamState {
broadcaster,
stop_token: CancellationToken::new(),
burst_buffer: Arc::new(Mutex::new(BurstBuffer::new(burst_buffer_size_in_bytes))),
task_handles: Mutex::new(Vec::new()),
}
}
@@ -158,6 +161,7 @@ impl SharedStreamState {
let (client_tx, client_rx) = mpsc::channel(self.buf_size);
let mut broadcast_rx = self.broadcaster.subscribe();
let cancel_token = CancellationToken::new();
{
let mut subs = self.subscribers.write().await;
subs.insert(*addr, cancel_token.clone());
@@ -169,8 +173,13 @@ impl SharedStreamState {
let burst_buffer_for_log = Arc::clone(&self.burst_buffer);
let yield_counter = YIELD_COUNTER;
// If a client stops streaming (for example presses
let timeout_duration = Duration::from_secs(300); // 5 minutes
let mut last_active = Instant::now();
let address = *addr;
tokio::spawn(async move {
let handle = tokio::spawn(async move {
// initial burst buffer
let snapshot = {
let buffer = burst_buffer.lock().await;
buffer.snapshot()
@@ -180,47 +189,63 @@ impl SharedStreamState {
let mut loop_cnt = 0;
loop {
tokio::select! {
biased;
biased;
() = cancel_token.cancelled() => {
debug!("Client disconnected from shared stream: {address}");
// canceled
() = cancel_token.cancelled() => {
debug!("Client disconnected from shared stream: {address}");
break;
}
// timeout handling
() = sleep(Duration::from_secs(1)) => {
if last_active.elapsed() > timeout_duration {
debug!("Client timed out due to inactivity: {address}");
cancel_token.cancel();
break;
}
result = broadcast_rx.recv() => {
match result {
Ok(data) => {
if let Err(err) = client_tx.send(data).await {
debug!("Shared stream client send error: {address} {err}");
break;
}
loop_cnt += 1;
if loop_cnt >= yield_counter {
tokio::task::yield_now().await;
loop_cnt = 0;
}
}
// receive broadcast data
result = broadcast_rx.recv() => {
match result {
Ok(data) => {
// Wenn der Client pausiert, einfach skippen oder warten
if client_tx_clone.is_closed() {
continue;
}
Err(tokio::sync::broadcast::error::RecvError::Lagged(skipped)) => {
let buffered_bytes = {
let buffer = burst_buffer_for_log.lock().await;
buffer.current_bytes
};
warn!("Shared stream client lagged behind {address}. Skipped {skipped} messages (buffered {buffered_bytes} bytes, yield counter {yield_counter})");
loop_cnt += 1;
if loop_cnt >= yield_counter {
tokio::task::yield_now().await;
loop_cnt = 0;
}
if let Err(err) = client_tx.send(data).await {
debug!("Shared stream client send error: {address} {err}");
break;
}
loop_cnt += 1;
last_active = Instant::now();
if loop_cnt >= yield_counter {
tokio::task::yield_now().await;
loop_cnt = 0;
}
Err(_) => break,
}
Err(tokio::sync::broadcast::error::RecvError::Lagged(skipped)) => {
let buffered_bytes = {
let buffer = burst_buffer_for_log.lock().await;
buffer.current_bytes
};
warn!("Shared stream client lagged behind {address}. Skipped {skipped} messages (buffered {buffered_bytes} bytes, yield counter {yield_counter})");
}
Err(_) => break,
}
}
}
}
manager.release_connection(&address, false).await;
});
let provider = self.provider_guard.as_ref().and_then(|h| h.allocation.get_provider_name());
self.task_handles.lock().await.push(handle);
let provider = self.provider_guard.as_ref().and_then(|h| h.allocation.get_provider_name());
(convert_stream(ReceiverStream::new(client_rx).boxed()), provider)
}
@@ -325,34 +350,37 @@ impl SharedStreamManager {
self.get_shared_state(stream_url).await.map(|s| s.headers.clone())
}
async fn unregister(&self, stream_url: &str, send_stop_signal: bool)
{
let shared_state = {
async fn unregister(&self, stream_url: &str, send_stop_signal: bool) {
let shared_state_opt = {
let mut shared_streams = self.shared_streams.write().await;
let shared_state = shared_streams.by_key.remove(stream_url);
{
let remove_keys: Vec<SocketAddr> = shared_streams.key_by_addr.iter()
.filter_map(|(addr, url)| if url == stream_url { Some(*addr) } else { None })
.collect();
for k in remove_keys {
shared_streams.key_by_addr.remove(&k);
}
let remove_keys: Vec<SocketAddr> = shared_streams.key_by_addr
.iter()
.filter_map(|(addr, url)| if url == stream_url { Some(*addr) } else { None })
.collect();
for k in remove_keys {
shared_streams.key_by_addr.remove(&k);
}
shared_state
shared_streams.by_key.remove(stream_url)
};
if let Some(shared_state) = shared_state {
if let Some(shared_state) = shared_state_opt {
let remaining = shared_state.subscribers.read().await.len();
debug_if_enabled!("Unregistering shared stream {} (remaining_subscribers={remaining}, send_stop_signal={send_stop_signal})",
sanitize_sensitive_info(stream_url));
sanitize_sensitive_info(stream_url));
for handle in shared_state.task_handles.lock().await.drain(..) {
handle.abort();
}
if let Some(provider_handle) = &shared_state.provider_guard {
self.provider_manager.release_handle(provider_handle).await;
}
if send_stop_signal {
if send_stop_signal || remaining == 0 {
trace_if_enabled!("Sending shared stream stop signal {}", sanitize_sensitive_info(stream_url));
let () = shared_state.stop_token.cancel();
shared_state.stop_token.cancel();
}
}
}
@@ -360,54 +388,65 @@ impl SharedStreamManager {
pub async fn release_connection(&self, addr: &SocketAddr, send_stop_signal: bool) {
let (stream_url, shared_state) = {
let shared_streams = self.shared_streams.read().await;
if let Some(stream_url) = shared_streams.key_by_addr.get(addr) {
(Some(stream_url.clone()), shared_streams.by_key.get(stream_url).cloned())
} else {
(None, None)
}
};
if let Some(stream_url) = shared_streams.key_by_addr.get(addr) {
(Some(stream_url.clone()), shared_streams.by_key.get(stream_url).cloned())
} else {
(None, None)
}
};
if let Some(state) = shared_state {
let (tx, is_empty, remaining) = {
let mut subs = state.subscribers.write().await;
let tx = subs.remove(addr);
let is_empty = subs.is_empty();
(
if send_stop_signal { tx } else { None },
is_empty,
subs.len(),
)
(tx, is_empty, subs.len())
};
debug!("Shared stream subscriber removed {addr}; remaining subscribers={remaining}");
if is_empty {
if let Some(url) = stream_url.as_ref() {
debug_if_enabled!("No subscribers remain for {} after removing {addr}", sanitize_sensitive_info(url)
);
debug_if_enabled!(
"No subscribers remain for {} after removing {addr}",
sanitize_sensitive_info(url)
);
self.unregister(url, send_stop_signal).await;
}
}
if let Some(client_stop_signal) = tx {
let () = client_stop_signal.cancel();
client_stop_signal.cancel();
}
}
}
async fn subscribe_stream(&self, stream_url: &str, addr: &SocketAddr, manager: Arc<SharedStreamManager>) -> Option<(BoxedProviderStream, Option<String>)> {
let mut shared_streams = self.shared_streams.write().await;
let shared_state_opt = shared_streams.by_key.get(stream_url).cloned();
match shared_state_opt {
Some(shared_state) => {
debug_if_enabled!("Responding to existing shared client stream {addr} {}", sanitize_sensitive_info(stream_url));
async fn subscribe_stream(
&self,
stream_url: &str,
addr: &SocketAddr,
manager: Arc<SharedStreamManager>,
) -> Option<(BoxedProviderStream, Option<String>)> {
let shared_state_opt = {
let shared_streams = self.shared_streams.read().await;
shared_streams.by_key.get(stream_url).cloned()
};
if let Some(shared_state) = shared_state_opt {
{
let mut shared_streams = self.shared_streams.write().await;
shared_streams.key_by_addr.insert(*addr, stream_url.to_owned());
Some(shared_state.subscribe(addr, manager).await)
}
None => None,
debug_if_enabled!("Responding to existing shared client stream {addr} {}",sanitize_sensitive_info(stream_url)
);
Some(shared_state.subscribe(addr, manager).await)
} else {
None
}
}
async fn register(&self, addr: &SocketAddr, stream_url: &str, shared_state: Arc<SharedStreamState>) {
let mut shared_streams = self.shared_streams.write().await;
shared_streams.by_key.insert(stream_url.to_string(), shared_state);
+35 -6
View File
@@ -1,14 +1,16 @@
use crate::model::{macros, EpgConfig};
use shared::error::{TuliproxError};
use shared::{check_input_connections, info_err, write_if_some};
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::utils::{get_base_url_from_str, get_credentials_from_url};
use shared::{check_input_credentials};
use shared::{check_input_connections, info_err, write_if_some};
use shared::check_input_credentials;
use std::collections::HashMap;
use std::fmt;
use std::path::PathBuf;
use url::Url;
use crate::utils::{get_csv_file_path};
#[allow(clippy::struct_excessive_bools)]
#[derive(Debug, Clone)]
@@ -101,6 +103,7 @@ pub struct ConfigInputAlias {
pub password: Option<String>,
pub priority: i16,
pub max_connections: u16,
pub exp_date: Option<i64>,
}
macros::from_impl!(ConfigInputAlias);
@@ -114,6 +117,7 @@ impl From<&ConfigInputAliasDto> for ConfigInputAlias {
password: dto.password.clone(),
priority: dto.priority,
max_connections: dto.max_connections,
exp_date: dto.exp_date,
}
}
}
@@ -136,6 +140,7 @@ pub struct ConfigInput {
pub max_connections: u16,
pub method: InputFetchMethod,
pub staged: Option<StagedInput>,
pub exp_date: Option<i64>,
pub t_batch_url: Option<String>,
}
@@ -150,6 +155,12 @@ impl ConfigInput {
return Err(info_err!("Staged input can only be from type m3u or xtream".to_owned()));
}
}
if is_input_expired(self.exp_date) {
warn!("Account {} expired for provider: {}", self.username.as_ref().map_or("?", |s| s.as_str()), self.name);
self.enabled = false;
}
Ok(batch_file_path)
}
@@ -180,10 +191,16 @@ impl ConfigInput {
InputType::Xtream
};
self.t_batch_url= Some(self.url.clone());
self.t_batch_url = Some(self.url.clone());
let file_path = get_csv_file_path(self.url.as_str()).ok();
if let Some(aliases) = self.aliases.as_mut() {
for alias in aliases.iter() {
if is_input_expired(alias.exp_date) {
warn!("Alias-Account {} expired for provider: {}", alias.username.as_ref().map_or("?", |s| s.as_str()), alias.name);
}
}
if !aliases.is_empty() {
let mut first = aliases.remove(0);
self.id = first.id;
@@ -223,6 +240,7 @@ impl ConfigInput {
max_connections: alias.max_connections,
method: self.method,
staged: None,
exp_date: None,
t_batch_url: None,
}
}
@@ -247,8 +265,9 @@ impl From<&ConfigInputDto> for ConfigInput {
priority: dto.priority,
max_connections: dto.max_connections,
method: dto.method,
t_batch_url: None,
exp_date: dto.exp_date,
staged: dto.staged.as_ref().map(StagedInput::from),
t_batch_url: None,
}
}
}
@@ -277,3 +296,13 @@ impl fmt::Display for ConfigInput {
Ok(())
}
}
pub fn is_input_expired(exp_date: Option<i64>) -> bool {
match exp_date {
Some(ts) => {
let now = Utc::now().timestamp();
ts <= now
}
None => false,
}
}
+3 -2
View File
@@ -230,6 +230,7 @@ pub fn get_attr_value(attr: &quick_xml::events::attributes::Attribute) -> Option
attr.unescape_value().ok().map(|v| v.to_string())
}
// This function filters a timeslot starting from yesterday.
#[allow(clippy::too_many_lines)]
async fn parse_xmltv_for_web_ui<R: AsyncRead + Send + Unpin>(reader: R) -> Result<EpgTv, TuliproxError> {
@@ -246,11 +247,11 @@ async fn parse_xmltv_for_web_ui<R: AsyncRead + Send + Unpin>(reader: R) -> Resul
// only 1 day old epg
let now = Utc::now();
let yesterday_start = Utc.with_ymd_and_hms(now.year(), now.month(), now.day(), 0, 0, 0).unwrap()
let yesterday_start = Utc.with_ymd_and_hms(now.year(), now.month(), now.day(), 0, 0, 0)
.single().expect("Current date at midnight should always be valid")
- chrono::Duration::days(1);
let threshold_ts = yesterday_start.timestamp();
loop {
match reader.read_event_into_async(&mut buf).await {
Ok(Event::Empty(e) | Event::Start(e)) => {
+7 -1
View File
@@ -2,12 +2,18 @@ use crate::model::{AppConfig, ProxyUserCredentials};
use crate::model::{ConfigTarget, XtreamTargetOutput};
use serde::{Deserialize, Deserializer, Serialize};
use serde_json::{Map, Value};
use shared::model::{xtream_const, PlaylistItem, XtreamPlaylistItem};
use shared::model::{xtream_const, PlaylistItem, ProxyUserStatus, XtreamPlaylistItem};
use shared::model::{ClusterFlags, PlaylistEntry, XtreamCluster};
use shared::utils::{deserialize_as_option_string, deserialize_as_string, deserialize_as_string_array, deserialize_number_from_string,
get_non_empty_str, opt_string_or_number_u32, string_default_on_null, string_or_number_f64, string_or_number_u32};
use std::iter::FromIterator;
#[derive(Debug, Default)]
pub struct XtreamLoginInfo {
pub status: Option<ProxyUserStatus>,
pub exp_date: Option<i64>,
}
#[derive(Deserialize, Default)]
pub struct XtreamCategory {
#[serde(deserialize_with = "deserialize_as_string")]
+19 -9
View File
@@ -7,20 +7,30 @@ use shared::error::{notify_err, TuliproxError};
use std::path::Path;
use tokio::io::AsyncWriteExt;
// Due to an error in quick_xml we cant write doc type through event. The quotes are escaped and the xml file is invalid.
//
// // XML Header
// writer.write_event_async(quick_xml::events::Event::Decl(quick_xml::events::BytesDecl::new("1.0", Some("utf-8"), None)))
// .await.map_err(|e| notify_err!(format!("failed to write XML header: {}", e)))?;
//
// // DOCTYPE
// writer.write_event_async(quick_xml::events::Event::DocType(quick_xml::events::BytesText::new(r#"tv SYSTEM "xmltv.dtd""#)))
// .await.map_err(|e| notify_err!(format!("failed to write doctype: {}", e)))?;
pub async fn epg_write_file(target: &ConfigTarget, epg: &Epg, path: &Path) -> Result<(), TuliproxError> {
let file = tokio::fs::File::create(path).await
.map_err(|e| notify_err!(format!("failed to create epg file: {}", e)))?;
let buf_writer = tokio::io::BufWriter::new(file);
let mut buf_writer = tokio::io::BufWriter::new(file);
// Work-Around BytesText DocType escape, see below
buf_writer.write_all(b"<?xml version=\"1.0\" encoding=\"utf-8\"?>\n").await
.map_err(|e| notify_err!(format!("failed to write XML header: {}", e)))?;
buf_writer.write_all(b"<!DOCTYPE tv SYSTEM \"xmltv.dtd\">\n").await
.map_err(|e| notify_err!(format!("failed to write doctype: {}", e)))?;
let mut writer = quick_xml::writer::Writer::new(buf_writer);
// XML Header
writer.write_event_async(quick_xml::events::Event::Decl(quick_xml::events::BytesDecl::new("1.0", Some("utf-8"), None)))
.await.map_err(|e| notify_err!(format!("failed to write XML header: {}", e)))?;
// DOCTYPE
writer.write_event_async(quick_xml::events::Event::DocType(quick_xml::events::BytesText::new(r#"tv SYSTEM "xmltv.dtd""#)))
.await.map_err(|e| notify_err!(format!("failed to write doctype: {}", e)))?;
// EPG Content
epg.write_to_async(&mut writer).await.map_err(|e| notify_err!(format!("failed to write epg: {}", e)))?;
+10 -16
View File
@@ -8,7 +8,7 @@ use crate::utils;
use crate::utils::json_write_documents_to_file;
use chrono::Local;
use log::error;
use shared::model::{PlaylistBouquetDto, ProxyType, ProxyUserStatus, TargetBouquetDto, TargetType, XtreamCluster};
use shared::model::{PlaylistBouquetDto, PlaylistClusterBouquetDto, ProxyType, ProxyUserStatus, TargetType, XtreamCluster};
use std::collections::{HashMap, HashSet};
use std::io::Error;
use std::path::{Path, PathBuf};
@@ -256,7 +256,6 @@ async fn save_xtream_user_bouquet_for_target(config: &Config, target_name: &str,
XtreamCluster::Series => user_get_series_bouquet_path(storage_path, TargetType::Xtream),
};
if let Some(bouquet_categories) = bouquet {
if let Some(xtream_categories) = xtream_get_playlist_categories(config, target_name, cluster).await {
let filtered: Vec<PlaylistXtreamCategory> = xtream_categories.iter().filter(|p| bouquet_categories.contains(&p.name)).cloned().collect();
@@ -294,27 +293,23 @@ async fn save_m3u_user_bouquet_for_target(storage_path: &Path, target: TargetTyp
Ok(())
}
async fn save_user_bouquet_for_target(config: &Config, target_name: &str, storage_path: &Path, target: TargetType, bouquet: &TargetBouquetDto) -> Result<(), Error> {
async fn save_user_bouquet_for_target(config: &Config, target_name: &str, storage_path: &Path, target: TargetType, bouquet: Option<&PlaylistClusterBouquetDto>) -> Result<(), Error> {
if target == TargetType::Xtream {
save_xtream_user_bouquet_for_target(config, target_name, storage_path, XtreamCluster::Live, bouquet.live.as_ref()).await?;
save_xtream_user_bouquet_for_target(config, target_name, storage_path, XtreamCluster::Video, bouquet.vod.as_ref()).await?;
save_xtream_user_bouquet_for_target(config, target_name, storage_path, XtreamCluster::Series, bouquet.series.as_ref()).await?;
save_xtream_user_bouquet_for_target(config, target_name, storage_path, XtreamCluster::Live, bouquet.and_then(|b| b.live.as_ref())).await?;
save_xtream_user_bouquet_for_target(config, target_name, storage_path, XtreamCluster::Video, bouquet.and_then(|b| b.vod.as_ref())).await?;
save_xtream_user_bouquet_for_target(config, target_name, storage_path, XtreamCluster::Series, bouquet.and_then(|b| b.series.as_ref())).await?;
} else {
save_m3u_user_bouquet_for_target(storage_path, target, XtreamCluster::Live, bouquet.live.as_ref()).await?;
save_m3u_user_bouquet_for_target(storage_path, target, XtreamCluster::Video, bouquet.vod.as_ref()).await?;
save_m3u_user_bouquet_for_target(storage_path, target, XtreamCluster::Series, bouquet.series.as_ref()).await?;
save_m3u_user_bouquet_for_target(storage_path, target, XtreamCluster::Live, bouquet.and_then(|b| b.live.as_ref())).await?;
save_m3u_user_bouquet_for_target(storage_path, target, XtreamCluster::Video, bouquet.and_then(|b| b.vod.as_ref())).await?;
save_m3u_user_bouquet_for_target(storage_path, target, XtreamCluster::Series, bouquet.and_then(|b| b.series.as_ref())).await?;
}
Ok(())
}
pub async fn save_user_bouquet(cfg: &Config, target_name: &str, username: &str, bouquet: &PlaylistBouquetDto) -> Result<(), Error> {
if let Some(storage_path) = ensure_user_storage_path(cfg, username) {
if let Some(xb) = &bouquet.xtream {
save_user_bouquet_for_target(cfg, target_name, &storage_path, TargetType::Xtream, xb).await?;
}
if let Some(mb) = &bouquet.m3u {
save_user_bouquet_for_target(cfg, target_name, &storage_path, TargetType::M3u, mb).await?;
}
save_user_bouquet_for_target(cfg, target_name, &storage_path, TargetType::Xtream, bouquet.xtream.as_ref()).await?;
save_user_bouquet_for_target(cfg, target_name, &storage_path, TargetType::M3u, bouquet.m3u.as_ref()).await?;
Ok(())
} else {
Err(Error::new(std::io::ErrorKind::NotFound, format!("User config path not found for user {username}")))
@@ -391,7 +386,6 @@ pub async fn user_get_bouquet_filter(config: &Config, username: &str, category_i
XtreamCluster::Series => user_get_series_bouquet(config, username, target).await,
};
match bouquet {
None => None,
Some(bouquet_categories) => {
+14 -5
View File
@@ -7,7 +7,7 @@ use std::io::{BufRead, Cursor, Error};
use std::path::PathBuf;
use url::Url;
use shared::model::{ConfigInputAliasDto, InputType};
use shared::utils::{get_credentials_from_url, trim_last_slash};
use shared::utils::{get_credentials_from_url, parse_timestamp, trim_last_slash};
use crate::utils::request::get_local_file_content;
const CSV_SEPARATOR: char = ';';
@@ -18,8 +18,9 @@ const FIELD_URL: &str = "url";
const FIELD_NAME: &str = "name";
const FIELD_USERNAME: &str = "username";
const FIELD_PASSWORD: &str = "password";
const FIELD_EXP_DATE: &str = "exp_date";
const FIELD_UNKNOWN: &str = "?";
const DEFAULT_COLUMNS: &[&str] = &[FIELD_URL, FIELD_MAX_CON, FIELD_PRIO, FIELD_NAME, FIELD_USERNAME, FIELD_PASSWORD];
const DEFAULT_COLUMNS: &[&str] = &[FIELD_URL, FIELD_MAX_CON, FIELD_PRIO, FIELD_NAME, FIELD_USERNAME, FIELD_PASSWORD, FIELD_EXP_DATE];
fn csv_assign_mandatory_fields(alias: &mut ConfigInputAliasDto, input_type: InputType) {
if !alias.url.is_empty() {
@@ -86,6 +87,12 @@ fn csv_assign_config_input_column(config_input: &mut ConfigInputAliasDto, header
FIELD_PASSWORD => {
config_input.password = Some(value.to_string());
}
FIELD_EXP_DATE => {
config_input.exp_date = parse_timestamp(value).unwrap_or_else(|e| {
error!("Failed to parse exp_date '{value}': {e}");
None
});
}
_ => {}
}
}
@@ -117,6 +124,7 @@ pub fn csv_read_inputs_from_reader(batch_input_type: InputType, reader: impl Buf
FIELD_NAME => FIELD_NAME,
FIELD_USERNAME => FIELD_USERNAME,
FIELD_PASSWORD => FIELD_PASSWORD,
FIELD_EXP_DATE => FIELD_EXP_DATE,
_ => {
error!("Field {s} is unsupported for csv input");
FIELD_UNKNOWN
@@ -135,6 +143,7 @@ pub fn csv_read_inputs_from_reader(batch_input_type: InputType, reader: impl Buf
password: None,
priority: 0,
max_connections: 1,
exp_date: None,
};
let columns: Vec<&str> = line.split(CSV_SEPARATOR).collect();
@@ -193,9 +202,9 @@ http://hd.providerline.com/get.php?username=user4&password=user4&type=m3u_plus;i
";
const XTREAM_BATCH: &str = r"
#name;username;password;url;max_connections
input_1;de566567;de2345f43g5;http://provider_1.tv:80;1
input_2;de566567;de2345f43g5;http://provider_2.tv:8080;1
#name;username;password;url;max_connections;exp_date
input_1;de566567;de2345f43g5;http://provider_1.tv:80;1;2028-11-23 13:12:34
input_2;de566567;de2345f43g5;http://provider_2.tv:8080;1;2028-12-23 13:12:34
";
#[test]
+99 -79
View File
@@ -1,20 +1,21 @@
use crate::api::model::AppState;
use crate::messaging::send_message;
use crate::model::{Config, ConfigInput, ConfigTarget};
use crate::model::{is_input_expired, Config, ConfigInput, ConfigTarget, XtreamLoginInfo};
use crate::model::{InputSource, ProxyUserCredentials};
use crate::processing::parser::xtream;
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::utils::request;
use chrono::{DateTime, Utc};
use log::{error, info, warn};
use shared::error::{str_to_io_error, TuliproxError};
use shared::model::{MsgKind, PlaylistEntry, PlaylistGroup, ProxyUserStatus, XtreamCluster, XtreamPlaylistItem};
use shared::utils::{extract_extension_from_url, get_i64_from_serde_value, get_string_from_serde_value};
// use std::cmp::Ordering;
use std::io::Error;
use std::str::FromStr;
use std::sync::Arc;
use std::time::{SystemTime, UNIX_EPOCH};
const THREE_DAYS_IN_SECS: i64 = 3 * 24 * 60 * 60;
#[inline]
pub fn get_xtream_stream_url_base(url: &str, username: &str, password: &str) -> String {
@@ -134,7 +135,7 @@ 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)];
async fn xtream_login(cfg: &Config, client: &Arc<reqwest::Client>, input: &InputSource, username: &str) -> Result<(), TuliproxError> {
async fn xtream_login(cfg: &Config, client: &Arc<reqwest::Client>, input: &InputSource, username: &str) -> Result<Option<XtreamLoginInfo>, TuliproxError> {
let content = if let Ok(content) = request::get_input_json_content(Arc::clone(client), None, input, None).await {
content
} else {
@@ -148,47 +149,57 @@ async fn xtream_login(cfg: &Config, client: &Arc<reqwest::Client>, input: &Input
}
};
match content.get("user_info") {
None => {}
Some(value) => {
if let Some(status_value) = value.get("status") {
if let Some(status) = 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:?}")).await;
}
}
}
}
if let Some(status_value) = value.get("exp_date") {
if let Some(expiration_timestamp) = 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 {
if let Some(datetime) = chrono::DateTime::from_timestamp(expiration_timestamp, 0) {
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}")).await;
}
}
} 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")).await;
}
}
let mut login_info = XtreamLoginInfo {
status: None,
exp_date: None,
};
if let Some(user_info) = content.get("user_info") {
if let Some(status_value) = user_info.get("status") {
if let Some(status) = get_string_from_serde_value(status_value) {
if let Ok(cur_status) = ProxyUserStatus::from_str(&status) {
login_info.status = Some(cur_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:?}")).await;
}
}
}
}
if let Some(exp_value) = user_info.get("exp_date") {
if let Some(expiration_timestamp) = get_i64_from_serde_value(exp_value) {
login_info.exp_date = Some(expiration_timestamp);
notify_account_expire(login_info.exp_date, cfg, client, username).await;
}
}
}
Ok(())
if login_info.exp_date.is_none() && login_info.status.is_none() {
Ok(None)
} else {
Ok(Some(login_info))
}
}
pub async fn notify_account_expire(exp_date: Option<i64>, cfg: &Config, client: &Arc<reqwest::Client>, username: &str) {
if let Some(expiration_timestamp) = exp_date {
let now_secs = Utc::now().timestamp(); // UTC-Time
if expiration_timestamp > now_secs {
let time_left = expiration_timestamp - now_secs;
if time_left < THREE_DAYS_IN_SECS {
if let Some(datetime) = DateTime::<Utc>::from_timestamp(expiration_timestamp, 0) {
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}")).await;
}
}
} 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")).await;
}
}
}
pub async fn get_xtream_playlist(cfg: &Arc<Config>, client: &Arc<reqwest::Client>, input: &Arc<ConfigInput>, working_dir: &str) -> (Vec<PlaylistGroup>, Vec<TuliproxError>) {
@@ -204,7 +215,9 @@ pub async fn get_xtream_playlist(cfg: &Arc<Config>, client: &Arc<reqwest::Client
let base_url = get_xtream_stream_url_base(&input_source.url, username, password);
let input_source_login = input_source.with_url(base_url.clone());
check_alias_user_state(cfg, client, input);
check_alias_user_state(cfg, client, input).await;
if let Err(err) = xtream_login(cfg, client, &input_source_login, username).await {
error!("Could not log in with xtream user {username} for provider {}. {err}", input.name);
return (Vec::with_capacity(0), vec![err]);
@@ -252,46 +265,53 @@ pub async fn get_xtream_playlist(cfg: &Arc<Config>, client: &Arc<reqwest::Client
(playlist_groups, errors)
}
/// Spawns a background task to check alias user states asynchronously.
/// This is fire-and-forget; errors are logged but not returned.
fn check_alias_user_state(cfg: &Arc<Config>, client: &Arc<reqwest::Client>, input: &Arc<ConfigInput>) {
let cfg = Arc::clone(cfg);
let client = Arc::clone(client);
let input = Arc::clone(input);
tokio::spawn(async move {
if let Some(aliases) = input.aliases.as_ref() {
for alias in aliases {
// Random wait time 5–20 seconds to avoid provider block
let delay = u64::from(fastrand::u32(5..=20));
tokio::time::sleep(tokio::time::Duration::from_secs(delay)).await;
if let (Some(username), Some(password)) =
(alias.username.as_ref(), alias.password.as_ref())
{
let mut input_source: InputSource = input.as_ref().into();
input_source.username.clone_from(&alias.username);
input_source.password.clone_from(&alias.password);
input_source.url.clone_from(&alias.url);
let base_url = get_xtream_stream_url_base(
&input_source.url,
username,
password,
);
let input_source_login = input_source.with_url(base_url.clone());
if let Err(err) =
xtream_login(&cfg, &client, &input_source_login, username).await
{
error!(
"Could not log in with xtream user {} for provider {}. {err}",
username,
alias.name
);
}
}
async fn check_alias_user_state(cfg: &Arc<Config>, client: &Arc<reqwest::Client>, input: &Arc<ConfigInput>) {
if let Some(aliases) = input.aliases.as_ref() {
for alias in aliases {
if is_input_expired(alias.exp_date) {
notify_account_expire(alias.exp_date, cfg, client, alias.username.as_ref().map_or("", |s| s.as_str())).await;
}
}
});
}
// TODO figure out how and when to call it to avoid provider bans. Possible reason for provider ban is to avoid brute force attacks.
//
// let cfg = Arc::clone(cfg);
// let client = Arc::clone(client);
// let input = Arc::clone(input);
//
// tokio::spawn(async move {
// for alias in &aliases {
// // Random wait time 60–180 seconds to avoid provider block
// let delay = u64::from(fastrand::u32(60..=180));
// tokio::time::sleep(tokio::time::Duration::from_secs(delay)).await;
//
// if let (Some(username), Some(password)) =
// (alias.username.as_ref(), alias.password.as_ref())
// {
// let mut input_source: InputSource = input.as_ref().into();
// input_source.username.clone_from(&alias.username);
// input_source.password.clone_from(&alias.password);
// input_source.url.clone_from(&alias.url);
// let base_url = get_xtream_stream_url_base(
// &input_source.url,
// username,
// password,
// );
// let input_source_login = input_source.with_url(base_url.clone());
//
// match xtream_login(&cfg, &client, &input_source_login, username).await {
// Ok(Some(xtream_login_info)) => {
// // TODO need to update the alias
//
// }
// Ok(None) => error!("Could log in with xtream user {} for provider {}. But could not extract account info", username, alias.name),
// Err(err) => error!("Could not log in with xtream user {} for provider {}. {err}",username,alias.name),
// }
// }
// }
// });
}
pub fn create_vod_info_from_item(target: &ConfigTarget, user: &ProxyUserCredentials, pli: &XtreamPlaylistItem, last_updated: i64) -> String {
@@ -315,4 +335,4 @@ pub fn create_vod_info_from_item(target: &ConfigTarget, user: &ProxyUserCredenti
"stream_id": {stream_id}
}}
}}"#)
}
}
+2 -2
View File
@@ -1,10 +1,10 @@
[package]
name = "frontend"
version = "3.2.10"
version = "3.2.11"
edition = "2021"
[dependencies]
shared = { version = "3.2.10", path = "../shared" }
shared = { version = "3.2.11", path = "../shared" }
chrono = "0"
yew = "0.21"
yew-router = "0.18"
+6 -1
View File
@@ -8,6 +8,7 @@
"LIVE_SHORT": "L",
"VOD_SHORT": "V",
"SERIES_SHORT": "S",
"MOVIE": "Movie",
"SAVE": "Save",
"SUBMIT": "Submit",
"OK": "Ok",
@@ -327,7 +328,7 @@
"COPY_CREDENTIALS": "Copy Credentials"
},
"TITLE": {
"USER_BOUQUET_EDITOR": "User group editor"
"USER_BOUQUET_EDITOR": "Playlist Category Selection"
},
"MESSAGES": {
"NO_CONTENT": "No content",
@@ -341,6 +342,7 @@
"SCHEDULE_EXISTS": "Schedule already exists",
"CLIPBOARD_NOT_SUPPORTED": "Clipboard not supported.\nYour browser or current context does not allow clipboard access.\nPlease use HTTPS or localhost.",
"FAILED_TO_KICK_USER_STREAM": "Failed to kick user stream",
"FAILED_TO_RETRIEVE_WEBPLAYER_URL": "Failed to retrieve webplayer URL",
"DOWNLOAD": {
"SUCCESS": "Successfully downloaded",
"FAIL": "Failed to download!",
@@ -354,6 +356,9 @@
"GEOIP": {
"SUCCESS": "Successfully downloaded Geo-IP db",
"FAIL": "Failed to download Geo-IP db!"
},
"USER_BOUQUET": {
"FAIL": "Failed to download user bouquets!"
}
},
"LOGIN": {
+7 -1
View File
@@ -201,7 +201,13 @@
"keys": [
"Checked"
],
"path": "M 4.2222224,2 C 3,2 2,3.000021 2,4.2222362 V 19.777764 C 2,20.999979 3,22 4.2222224,22 H 19.777778 C 20.999989,22 22,20.999979 22,19.777764 V 4.2222362 C 22,3.000021 20.999989,2 19.777778,2 Z m 0,2.2222362 H 19.777778 V 19.777764 H 4.2222224 Z M 17.288622,7.1063517 10.441833,13.953134 6.7113668,10.233533 5.2465223,11.698352 10.441833,16.89369 18.753467,8.5820472 Z"
"path": "M12 7c-2.76 0-5 2.24-5 5s2.24 5 5 5 5-2.24 5-5-2.24-5-5-5zm0-5C6.48 2 2 6.48 2 12s4.48 10 10 10 10-4.48 10-10S17.52 2 12 2zm0 18c-4.42 0-8-3.58-8-8s3.58-8 8-8 8 3.58 8 8-3.58 8-8 8z"
},
{
"keys": [
"Unchecked"
],
"path": "M12 2C6.48 2 2 6.48 2 12s4.48 10 10 10 10-4.48 10-10S17.52 2 12 2zm0 18c-4.42 0-8-3.58-8-8s3.58-8 8-8 8 3.58 8 8-3.58 8-8 8z"
},
{
"keys": [
+2
View File
@@ -17,6 +17,8 @@ $form-grid-cell-width-small-screen: 200px;
$form-field-min-width-small-screen: 160px;
$form-max-grid-cells-small-screen: 2;
$playlist-categories-grid-cell-width: 400px;
@mixin fixed-width($width) {
width: $width;
min-width: $width;
@@ -1,3 +1,136 @@
div {
color: var(--text-color);
@use "../../../size" as size;
.tp__api-user-playlist {
display: flex;
flex-flow: column;
gap: var(--gap-default);
box-sizing: border-box;
overflow: hidden;
padding: var(--padding-small);
width: 100%;
&__loading {
transform: translateY(-10px);
}
&__header {
flex-flow: row wrap;
gap: var(--gap-default);
&-toolbar {
display: flex;
flex-flow: row wrap;
gap: var(--gap-default);
box-sizing: border-box;
}
}
&__content {
display: flex;
flex-flow: column;
gap: var(--gap-default);
box-sizing: border-box;
overflow: hidden;
width: 100%;
&-toolbar {
display: flex;
flex-flow: row;
box-sizing: border-box;
width: 100%;
justify-content: space-between;
}
&-panels {
display: flex;
flex-flow: column;
width: 100%;
gap: var(--gap-default);
box-sizing: border-box;
overflow: hidden;
}
}
}
.tp__api-user-target-playlist {
display: flex;
flex-flow: column;
width: 100%;
gap: var(--gap-default);
box-sizing: border-box;
overflow: hidden;
&__body {
display: flex;
flex-flow: column;
width: 100%;
gap: var(--gap-default);
padding: var(--padding-default) 0;
box-sizing: border-box;
overflow: auto;
.tp__collapse-panel__header {
font-size: 1.6rem;
}
> :nth-child(1) {
border-left: 2px solid var(--output-m3u-color);
}
> :nth-child(2) {
border-left: 2px solid var(--output-xtream-color);
}
> :nth-child(3) {
border-left: 2px solid var(--output-hdhomerun-color);
}
}
&__categories {
display: grid;
grid-template-columns: repeat(auto-fit, minmax(size.$playlist-categories-grid-cell-width, 1fr));
gap: var(--gap-large);
padding: var(--padding-default) 0;
box-sizing: border-box;
width: 100%;
&-category.selected {
background-color: var(--text-button-active-background-color);
color: var(--text-button-active-color);
fill: var(--text-button-active-color);
}
&-category {
display: flex;
justify-content: flex-start;
align-items: center;
background-color: var(--text-button-background-color);
color: var(--text-button-color);
fill: var(--text-button-color);
border: 1px solid var(--text-button-border-color);
border-radius: var(--border-radius);
cursor: pointer;
box-sizing: border-box;
max-height: 3rem;
min-height: 2.5rem;
gap: var(--gap-default);
padding: 0 var(--padding-default);
overflow: hidden;
text-overflow: ellipsis;
svg, img {
pointer-events: none;
height: 1.2rem;
width: 1.2rem;
}
&:hover {
background-color: var(--text-button-hover-background-color);
color: var(--text-button-hover-color);
}
}
}
}
@@ -5,7 +5,7 @@ use crate::app::components::theme::Theme;
use crate::hooks::use_service_context;
use crate::provider::DialogProvider;
use yew::use_state;
use crate::app::components::api_user::playlist::ApiUserPlaylist;
#[function_component]
pub fn ApiUserView() -> Html {
@@ -50,7 +50,7 @@ pub fn ApiUserView() -> Html {
</div>
</div>
<div class="tp__app-main__body">
{" TODO "}
<ApiUserPlaylist />
</div>
</div>
</div>
@@ -1,3 +1,5 @@
mod api_user_view;
mod playlist;
mod target_playlist;
pub use api_user_view::*;
@@ -0,0 +1,226 @@
use crate::app::components::api_user::target_playlist::{BouquetSelection, UserTargetPlaylist};
use crate::app::components::{Panel, RadioButtonGroup, TextButton};
use crate::hooks::use_service_context;
use crate::model::{BusyStatus, EventMessage};
use shared::error::TuliproxError;
use shared::info_err;
use shared::model::{PlaylistBouquetDto, PlaylistCategoriesDto, PlaylistClusterBouquetDto};
use std::cell::RefCell;
use std::collections::HashMap;
use std::fmt;
use std::rc::Rc;
use std::str::FromStr;
use yew::prelude::*;
use yew_i18n::use_translation;
#[derive(Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
enum ApiUserPlaylistPage {
Xtream,
M3u,
}
impl FromStr for ApiUserPlaylistPage {
type Err = TuliproxError;
fn from_str(s: &str) -> Result<Self, TuliproxError> {
match s.to_lowercase().as_str() {
"xtream" => Ok(ApiUserPlaylistPage::Xtream),
"m3u" => Ok(ApiUserPlaylistPage::M3u),
_ => Err(info_err!(format!("Unknown api user playlist type: {s}"))),
}
}
}
impl fmt::Display for ApiUserPlaylistPage {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
let s = match self {
ApiUserPlaylistPage::Xtream => "xtream",
ApiUserPlaylistPage::M3u => "m3u",
};
write!(f, "{s}")
}
}
fn to_playlist_cluster(count: (usize, usize, usize), bouquet: Option<&Rc<RefCell<BouquetSelection>>>) -> Option<PlaylistClusterBouquetDto> {
if let Some(bouq) = bouquet {
let selections = bouq.borrow();
let selected_vec = |map: &HashMap<String, bool>| {
let v: Vec<String> = map.iter()
.filter(|(_, &selected)| selected)
.map(|(c, _)| c.clone())
.collect();
if v.is_empty() { None } else { Some(v) }
};
let live = selected_vec(&selections.live).filter(|v| v.len() != count.0);
let vod = selected_vec(&selections.vod).filter(|v| v.len() != count.1);
let series = selected_vec(&selections.series).filter(|v| v.len() != count.2);
// if all three are None, return None
if live.is_none() && vod.is_none() && series.is_none() {
None
} else {
Some(PlaylistClusterBouquetDto { live, vod, series })
}
} else {
None
}
}
#[function_component]
pub fn ApiUserPlaylist() -> Html {
let translate = use_translation();
let service_ctx = use_service_context();
let categories = use_state(|| None as Option<Rc<PlaylistCategoriesDto>>);
let bouquets = use_state(|| None as Option<Rc<PlaylistBouquetDto>>);
let active_tab = use_state(|| ApiUserPlaylistPage::Xtream);
let playlist_types = use_memo((), |_| {
[ApiUserPlaylistPage::Xtream, ApiUserPlaylistPage::M3u].iter().map(ToString::to_string).collect::<Vec<String>>()
});
// Selection reference
let selections = use_mut_ref(|| {
HashMap::<ApiUserPlaylistPage, Rc<RefCell<BouquetSelection>>>::new()
});
let handle_tab_select = {
let active_tab_clone = active_tab.clone();
Callback::from(move |page_selection: Rc<Vec<String>>| {
if let Some(page_selection_str) = page_selection.first() {
if let Ok(page) = ApiUserPlaylistPage::from_str(page_selection_str) {
active_tab_clone.set(page)
}
}
})
};
{
// ----- Load data on mount -----
let categories = categories.clone();
let bouquets = bouquets.clone();
let services = service_ctx.clone();
let translate = translate.clone();
use_effect_with((), move |_| {
wasm_bindgen_futures::spawn_local(async move {
services.event.broadcast(EventMessage::Busy(BusyStatus::Show));
let result = (services.user_api.get_playlist_bouquet().await, services.user_api.get_playlist_categories().await);
match result {
(Ok(bouquet), Ok(cats)) => {
bouquets.set(bouquet.clone());
categories.set(cats.clone());
}
(Err(e1), Err(e2)) => {
log::error!("Failed to load bouquet: {e1:?}, categories: {e2:?}");
services.toastr.error(translate.t("MESSAGES.DOWNLOAD.USER_BOUQUET.FAIL"));
}
(Err(e), _) | (_, Err(e)) => {
log::error!("Failed to load user data: {e:?}");
services.toastr.error(translate.t("MESSAGES.DOWNLOAD.USER_BOUQUET.FAIL"));
}
}
services.event.broadcast(EventMessage::Busy(BusyStatus::Hide));
});
|| {}
});
}
// ----- Save handler -----
let on_save = {
let selections = selections.clone();
let services = service_ctx.clone();
let translate = translate.clone();
let categories = categories.clone();
Callback::from(move |_| {
let selections = selections.clone();
let services = services.clone();
let translate = translate.clone();
let categories_xtream_count = categories.as_ref().and_then(|plc| plc.xtream.as_ref().map(|x|
(x.live.as_ref().map(|v| v.len()).unwrap_or(0),
x.vod.as_ref().map(|v| v.len()).unwrap_or(0),
x.series.as_ref().map(|v| v.len()).unwrap_or(0))
)).unwrap_or((0, 0, 0));
let categories_m3u_count = categories.as_ref().and_then(|plc| plc.m3u.as_ref().map(|x|
(x.live.as_ref().map(|v| v.len()).unwrap_or(0),
x.vod.as_ref().map(|v| v.len()).unwrap_or(0),
x.series.as_ref().map(|v| v.len()).unwrap_or(0))
)).unwrap_or((0, 0, 0));
wasm_bindgen_futures::spawn_local(async move {
services.event.broadcast(EventMessage::Busy(BusyStatus::Show));
let result = {
let selects = selections.borrow();
PlaylistBouquetDto {
xtream: to_playlist_cluster(categories_xtream_count, selects.get(&ApiUserPlaylistPage::Xtream)),
m3u: to_playlist_cluster(categories_m3u_count, selects.get(&ApiUserPlaylistPage::M3u)),
}
};
match services.user_api.save_playlist_bouquet(&result).await {
Ok(()) => services.toastr.success(translate.t("MESSAGES.SAVE.BOUQUET.SUCCESS")),
Err(_) => services.toastr.error(translate.t("MESSAGES.SAVE.BOUQUET.FAIL")),
}
services.event.broadcast(EventMessage::Busy(BusyStatus::Hide));
});
})
};
let handle_m3u_change = {
let selections = selections.clone();
Callback::from(move |selection: Rc<RefCell<BouquetSelection>>| {
selections.borrow_mut().insert(ApiUserPlaylistPage::M3u, selection);
})
};
let handle_xtream_change = {
let selections = selections.clone();
Callback::from(move |selection: Rc<RefCell<BouquetSelection>>| {
selections.borrow_mut().insert(ApiUserPlaylistPage::Xtream, selection);
})
};
html! {
<div class="tp__api-user-playlist">
<div class="tp__api-user-playlist__header tp__list-list__header">
<h1>{translate.t("TITLE.USER_BOUQUET_EDITOR") }</h1>
<div class="tp__userlist-list__header-toolbar">
<TextButton class="primary" name="save"
icon="Save"
title={ translate.t("LABEL.SAVE")}
onclick={on_save}></TextButton>
</div>
</div>
<div class="tp__api-user-playlist__content">
<div class="user-playlist__content-toolbar">
<RadioButtonGroup options={playlist_types.clone()}
selected={Rc::new(vec![(*active_tab).to_string()])}
on_select={handle_tab_select} />
</div>
<div class="tp__api-user-playlist__content-panels">
<Panel value={ApiUserPlaylistPage::Xtream.to_string()} active={active_tab.to_string()}>
<UserTargetPlaylist
categories={categories.as_ref().and_then(|c| c.xtream.clone())}
bouquet={bouquets.as_ref().and_then(|b| b.xtream.as_ref().cloned())}
on_change={handle_xtream_change.clone()}
></UserTargetPlaylist>
</Panel>
<Panel value={ApiUserPlaylistPage::M3u.to_string()} active={active_tab.to_string()}>
<UserTargetPlaylist
categories={categories.as_ref().and_then(|c| c.m3u.clone())}
bouquet={bouquets.as_ref().and_then(|b| b.m3u.as_ref().cloned())}
on_change={handle_m3u_change.clone()}
></UserTargetPlaylist>
</Panel>
</div>
</div>
</div>
}
}
@@ -0,0 +1,177 @@
use std::cell::RefCell;
use crate::app::components::{AppIcon, Card, CollapsePanel};
use shared::model::{PlaylistClusterBouquetDto, PlaylistClusterCategoriesDto, XtreamCluster};
use std::collections::HashMap;
use std::rc::Rc;
use std::str::FromStr;
use wasm_bindgen::JsCast;
use yew::prelude::*;
use yew_i18n::use_translation;
use crate::html_if;
fn normalize(s: &str) -> String {
let cleaned: String = s
.chars()
.filter(|c| c.is_alphanumeric() || c.is_whitespace())
.collect();
cleaned.trim().to_lowercase()
}
fn sort_opt_vec(v: &mut Option<Vec<String>>) {
if let Some(ref mut inner) = v {
inner.sort_by_key(|a| normalize(a));
}
}
macro_rules! create_selection {
($bouquet:expr, $categories:expr, $selections:expr, $field: ident) => {
if let Some(selects) = $bouquet.$field.as_ref() {
for b in selects {
$selections.$field.insert(b.clone(), true);
}
} else {
if let Some(cats) = $categories.$field.as_ref() {
for c in cats {
$selections.$field.insert(c.clone(), true);
}
}
}
};
}
#[derive(Clone, PartialEq, Default)]
pub struct BouquetSelection {
pub live: HashMap<String, bool>,
pub vod: HashMap<String, bool>,
pub series: HashMap<String, bool>,
}
#[derive(Properties, PartialEq)]
pub struct UserTargetPlaylistProps {
pub categories: Option<PlaylistClusterCategoriesDto>,
pub bouquet: Option<PlaylistClusterBouquetDto>,
pub on_change: Callback<Rc<RefCell<BouquetSelection>>>,
}
#[function_component]
pub fn UserTargetPlaylist(props: &UserTargetPlaylistProps) -> Html {
let translate = use_translation();
let bouquet_selection = use_mut_ref(BouquetSelection::default);
let playlist_categories = use_state(PlaylistClusterCategoriesDto::default);
let force_update = use_state(|| 0);
{
let bouquet_selection = bouquet_selection.clone();
let playlist_categories = playlist_categories.clone();
let in_cats = props.categories.clone();
let in_bouquet = props.bouquet.clone();
let force_update = force_update.clone();
use_effect_with((in_cats, in_bouquet), move |(maybe_categories, maybe_bouquet)| {
let mut selections = BouquetSelection::default();
if let Some(categories) = maybe_categories.as_ref() {
if let Some(bouquet) = maybe_bouquet.as_ref() {
create_selection!(bouquet, categories, selections, live);
create_selection!(bouquet, categories, selections, vod);
create_selection!(bouquet, categories, selections, series);
} else {
if let Some(cats) = categories.live.as_ref() {
for c in cats {
selections.live.insert(c.clone(), true);
}
}
if let Some(cats) = categories.vod.as_ref() {
for c in cats {
selections.vod.insert(c.clone(), true);
}
}
if let Some(cats) = categories.series.as_ref() {
for c in cats {
selections.series.insert(c.clone(), true);
}
}
}
*bouquet_selection.borrow_mut() = selections;
let mut new_categories = categories.clone();
sort_opt_vec(&mut new_categories.live);
sort_opt_vec(&mut new_categories.vod);
sort_opt_vec(&mut new_categories.series);
playlist_categories.set(new_categories);
force_update.set(*force_update + 1);
}
});
}
let handle_category_click = {
let on_change = props.on_change.clone();
let bouquet_selection = bouquet_selection.clone();
let force_update = force_update.clone();
Callback::from(move |e: MouseEvent| {
e.prevent_default();
if let Some(target) = e.target() {
if let Ok(element) = target.dyn_into::<web_sys::Element>() {
if let Some(cluster) = element.get_attribute("data-cluster") {
if let Ok(cluster) = XtreamCluster::from_str(cluster.as_str()) {
if let Some(category) = element.get_attribute("data-category") {
let mut selections = bouquet_selection.borrow_mut();
match cluster {
XtreamCluster::Live => {
let selected = *selections.live.get(&category).unwrap_or(&false);
selections.live.insert(category, !selected);
}
XtreamCluster::Video => {
let selected = *selections.vod.get(&category).unwrap_or(&false);
selections.vod.insert(category, !selected);
}
XtreamCluster::Series => {
let selected = *selections.series.get(&category).unwrap_or(&false);
selections.series.insert(category, !selected);
}
}
on_change.emit(bouquet_selection.clone());
force_update.set(*force_update + 1);
}
}
}
}
}
})
};
let render_category_cluster = |cluster: XtreamCluster, cats: Option<&Vec<String>>, selections: &HashMap<String, bool>| {
if let Some(c) = cats {
html_if!(!c.is_empty(), {
<Card>
<CollapsePanel title={translate.t( match cluster {
XtreamCluster::Live => "LABEL.LIVE",
XtreamCluster::Video => "LABEL.MOVIE",
XtreamCluster::Series => "LABEL.SERIES"
})}>
<div class="tp__api-user-target-playlist__categories">
{ for c.iter().map(|cat| {
let selected = *selections.get(cat).unwrap_or(&false);
html! {
<div key={cat.clone()} data-cluster={cluster.to_string()} data-category={cat.clone()} class={classes!("tp__api-user-target-playlist__categories-category", if selected {"selected"} else {""})}
onclick={handle_category_click.clone()}>
<AppIcon name={if selected {"Checked"} else {"Unchecked"}}/> { &cat }
</div>
}})}
</div>
</CollapsePanel>
</Card>
})
} else {
html! {}
}
};
let selections = &*bouquet_selection.borrow();
html! {
<div class={"tp__api-user-target-playlist"}>
<div class="tp__api-user-target-playlist__body">
{ render_category_cluster(XtreamCluster::Live, playlist_categories.live.as_ref(), &selections.live) }
{ render_category_cluster(XtreamCluster::Video, playlist_categories.vod.as_ref(), &selections.vod) }
{ render_category_cluster(XtreamCluster::Series, playlist_categories.series.as_ref(), &selections.series) }
</div>
</div>
}
}
@@ -14,7 +14,9 @@ use std::rc::Rc;
use std::str::FromStr;
use wasm_bindgen::JsCast;
use web_sys::Element;
use yew::platform::spawn_local;
use yew::prelude::*;
use yew_hooks::use_clipboard;
use yew_i18n::use_translation;
const LIVE: &str = "Live";
@@ -24,6 +26,11 @@ const CATCHUP: &str = "Archive";
const HLS: &str = "HLS";
const DASH: &str = "DASH";
const KICK: &str = "kick";
const COPY_LINK_TULIPROX_VIRTUAL_ID: &str = "copy_link_tuliprox_virtual_id";
const COPY_LINK_TULIPROX_WEBPLAYER_URL: &str = "copy_link_tuliprox_webplayer_url";
const COPY_LINK_PROVIDER_URL: &str = "copy_link_provider_url";
const HEADERS: [&str; 12] = [
"EMPTY",
"USERNAME",
@@ -70,7 +77,8 @@ pub struct StreamsTableProps {
#[function_component]
pub fn StreamsTable(props: &StreamsTableProps) -> Html {
let translate = use_translation();
let services = use_service_context();
let service_ctx = use_service_context();
let clipboard = use_clipboard();
let config_ctx = use_context::<ConfigContext>().expect("Config context not found");
let popup_anchor_ref = use_state(|| None::<web_sys::Element>);
let popup_is_open = use_state(|| false);
@@ -214,22 +222,62 @@ pub fn StreamsTable(props: &StreamsTableProps) -> Html {
})
};
let copy_to_clipboard: Callback<String> = {
let clipboard = clipboard.clone();
let services = service_ctx.clone();
let translate = translate.clone();
Callback::from(move |text: String| {
if *clipboard.is_supported {
clipboard.write_text(text);
} else {
services.toastr.error(translate.t("MESSAGES.CLIPBOARD_NOT_SUPPORTED"));
}
})
};
let handle_menu_click = {
let popup_is_open_state = popup_is_open.clone();
let translate = translate.clone();
let services_ctx = services.clone();
let services = service_ctx.clone();
let selected_dto = selected_dto.clone();
let copy_to_clipboard = copy_to_clipboard.clone();
Callback::from(move |(name, _): (String, _)| {
if let Ok(action) = StreamsTableAction::from_str(&name) {
match action {
StreamsTableAction::Kick => {
if let Some(dto) = (*selected_dto).as_ref() {
if !services_ctx.websocket.send_message(ProtocolMessage::UserAction(UserCommand::Kick(dto.addr))) {
services_ctx.toastr.error(translate.t("MESSAGES.FAILED_TO_KICK_USER_STREAM"));
if !services.websocket.send_message(ProtocolMessage::UserAction(UserCommand::Kick(dto.addr))) {
services.toastr.error(translate.t("MESSAGES.FAILED_TO_KICK_USER_STREAM"));
}
}
}
StreamsTableAction::CopyLinkTuliproxVirtualId => {
if let Some(dto) = &*selected_dto {
copy_to_clipboard.emit(dto.channel.virtual_id.to_string());
}
}
StreamsTableAction::CopyLinkProviderUrl => {
if let Some(dto) = &*selected_dto {
copy_to_clipboard.emit(dto.channel.url.clone());
}
}
StreamsTableAction::CopyLinkTuliproxWebPlayerUrl => {
if let Some(dto) = &*selected_dto {
let target_id = dto.channel.target_id;
let virtual_id = dto.channel.virtual_id;
let cluster = dto.channel.cluster;
let services = services.clone();
let translate = translate.clone();
let copy_to_clipboard = copy_to_clipboard.clone();
spawn_local(async move {
if let Some(url) = services.playlist.get_playlist_webplayer_url(target_id, virtual_id, cluster).await {
copy_to_clipboard.emit(url);
} else {
services.toastr.error(translate.t("MESSAGES.FAILED_TO_RETRIEVE_WEBPLAYER_URL"));
}
});
}
}
}
}
popup_is_open_state.set(false);
@@ -245,6 +293,9 @@ pub fn StreamsTable(props: &StreamsTableProps) -> Html {
<Table::<StreamInfo> definition={definition.clone()} />
<PopupMenu is_open={*popup_is_open} anchor_ref={(*popup_anchor_ref).clone()} on_close={handle_popup_close}>
<MenuItem icon="Disconnect" name={StreamsTableAction::Kick.to_string()} label={translate.t("LABEL.KICK")} onclick={&handle_menu_click} class="tp__delete_action"></MenuItem>
<MenuItem icon="Clipboard" name={StreamsTableAction::CopyLinkTuliproxVirtualId.to_string()} label={translate.t("LABEL.COPY_LINK_TULIPROX_VIRTUAL_ID")} onclick={&handle_menu_click}></MenuItem>
<MenuItem icon="Clipboard" name={StreamsTableAction::CopyLinkTuliproxWebPlayerUrl.to_string()} label={translate.t("LABEL.COPY_LINK_TULIPROX_WEBPLAYER_URL")} onclick={&handle_menu_click}></MenuItem>
<MenuItem icon="Clipboard" name={StreamsTableAction::CopyLinkProviderUrl.to_string()} label={translate.t("LABEL.COPY_LINK_PROVIDER_URL")} onclick={&handle_menu_click}></MenuItem>
</PopupMenu>
</>
}
@@ -259,12 +310,18 @@ pub fn StreamsTable(props: &StreamsTableProps) -> Html {
#[derive(Debug, Clone, Eq, PartialEq)]
enum StreamsTableAction {
Kick,
CopyLinkTuliproxVirtualId,
CopyLinkTuliproxWebPlayerUrl,
CopyLinkProviderUrl,
}
impl Display for StreamsTableAction {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "{}", match self {
Self::Kick => "kick",
Self::Kick => KICK,
Self::CopyLinkTuliproxVirtualId => COPY_LINK_TULIPROX_VIRTUAL_ID,
Self::CopyLinkTuliproxWebPlayerUrl => COPY_LINK_TULIPROX_WEBPLAYER_URL,
Self::CopyLinkProviderUrl => COPY_LINK_PROVIDER_URL,
})
}
}
@@ -273,8 +330,14 @@ impl FromStr for StreamsTableAction {
type Err = TuliproxError;
fn from_str(s: &str) -> Result<Self, TuliproxError> {
if s.eq("kick") {
if s.eq(KICK) {
Ok(Self::Kick)
} else if s.eq(COPY_LINK_TULIPROX_VIRTUAL_ID) {
Ok(Self::CopyLinkTuliproxVirtualId)
} else if s.eq(COPY_LINK_TULIPROX_WEBPLAYER_URL) {
Ok(Self::CopyLinkTuliproxWebPlayerUrl)
} else if s.eq(COPY_LINK_PROVIDER_URL) {
Ok(Self::CopyLinkProviderUrl)
} else {
create_tuliprox_error_result!(TuliproxErrorKind::Info, "Unknown Stream Action: {}", s)
}
@@ -144,13 +144,16 @@ pub fn PlaylistExplorer() -> Html {
if let Some(dto) = &*selected_channel {
let copy_to_clipboard = copy_to_clipboard.clone();
let services = services.clone();
let dto = dto.clone();
let virtual_id = dto.virtual_id;
let cluster = dto.xtream_cluster.unwrap_or_default();
let translate_clone = translate_clone.clone();
let target_id = *target_id;
spawn_local(async move {
if let Some(url) = services.playlist.get_playlist_webplayer_url(target_id, &dto).await {
if let Some(url) = services.playlist.get_playlist_webplayer_url(target_id, virtual_id, cluster).await {
copy_to_clipboard.emit(url);
services.toastr.success(translate_clone.t("MESSAGES.PLAYLIST.WEBPLAYER_URL_COPY_TO_CLIPBOARD"));
} else {
services.toastr.error(translate_clone.t("MESSAGES.FAILED_TO_RETRIEVE_WEBPLAYER_URL"));
}
});
}
@@ -149,7 +149,7 @@ pub fn XtreamTargetOutputView(props: &XtreamTargetOutputViewProps) -> Html {
{render_output()}
</Panel>
<Panel value={OutputFormPage::Trakt.to_string()} active={view_visible.to_string()}>
{"TODO"}
{"TODO..."}
</Panel>
</div>
</div>
-13
View File
@@ -12,19 +12,6 @@ pub struct TabItem {
pub inactive_class: Option<String>,
}
// impl TabItem {
// pub fn new(id: String, title: String, icon: String, children: Html) -> Self {
// Self {
// id,
// title,
// icon,
// children,
// active_class: None,
// inactive_class: None,
// }
// }
// }
#[derive(Properties, Clone, PartialEq)]
pub struct TabSetProps {
pub tabs: Rc<Vec<TabItem>>,
+4 -1
View File
@@ -1,12 +1,13 @@
use std::rc::Rc;
use yew::prelude::*;
use crate::model::WebConfig;
use crate::services::{AuthService, ConfigService, EventService, PlaylistService, StatusService, StreamsService, ToastrService, UserService, WebSocketService};
use crate::services::{AuthService, ConfigService, EventService, PlaylistService, StatusService, StreamsService, ToastrService, UserApiService, UserService, WebSocketService};
pub struct Services {
pub auth: Rc<AuthService>,
pub config: Rc<ConfigService>,
pub user: Rc<UserService>,
pub user_api: Rc<UserApiService>,
pub status: Rc<StatusService>,
pub streams: Rc<StreamsService>,
pub event: Rc<EventService>,
@@ -25,6 +26,7 @@ impl Services {
let playlist = Rc::new(PlaylistService::new());
let toastr = Rc::new(ToastrService::new());
let user = Rc::new(UserService::new(Rc::clone(&event)));
let user_api = Rc::new(UserApiService::new());
let websocket = Rc::new(WebSocketService::new(Rc::clone(&status), Rc::clone(&event)));
Self {
auth,
@@ -34,6 +36,7 @@ impl Services {
event,
playlist,
user,
user_api,
toastr,
websocket
}
+2
View File
@@ -8,6 +8,7 @@ mod websocket_service;
mod toastr_service;
mod event_service;
mod user_service;
mod user_api_service;
mod streams_service;
pub use self::auth_service::*;
@@ -20,4 +21,5 @@ pub use self::websocket_service::*;
pub use self::toastr_service::*;
pub use self::event_service::*;
pub use self::user_service::*;
pub use self::user_api_service::*;
pub use self::streams_service::*;
+4 -4
View File
@@ -1,6 +1,6 @@
use crate::services::{get_base_href, request_post, ACCEPT_PREFER_BIN};
use log::error;
use shared::model::{CommonPlaylistItem, EpgTv, PlaylistCategoriesResponse, PlaylistEpgRequest, PlaylistRequest, UiPlaylistCategories, WebplayerUrlRequest};
use shared::model::{EpgTv, PlaylistCategoriesResponse, PlaylistEpgRequest, PlaylistRequest, UiPlaylistCategories, WebplayerUrlRequest, XtreamCluster};
use std::rc::Rc;
use shared::utils::{concat_path_leading_slash};
@@ -41,11 +41,11 @@ impl PlaylistService {
})
}
pub async fn get_playlist_webplayer_url(&self, target_id: u16, dto: &Rc<CommonPlaylistItem>) -> Option<String> {
pub async fn get_playlist_webplayer_url(&self, target_id: u16, virtual_id: u32, cluster: XtreamCluster) -> Option<String> {
let request = WebplayerUrlRequest {
target_id,
virtual_id: dto.virtual_id,
cluster: dto.xtream_cluster.unwrap_or_default(),
virtual_id,
cluster,
};
request_post::<&WebplayerUrlRequest, String>(&self.playlist_api_webplayer_url_path, &request, None, Some("text/plain".to_string())).await.unwrap_or_else(|err| {
error!("{err}");
+41
View File
@@ -0,0 +1,41 @@
use std::rc::Rc;
use log::error;
use shared::model::{PlaylistBouquetDto, PlaylistCategoriesDto};
use shared::utils::{concat_path_leading_slash};
use crate::error::Error;
use crate::services::{get_base_href, request_get, request_post};
#[derive(Debug, Default)]
pub struct UserApiService {
user_playlist_categories_path: String,
user_playlist_bouquet_path: String,
}
impl UserApiService {
pub fn new() -> Self {
let base_href = get_base_href();
Self {
user_playlist_categories_path: concat_path_leading_slash(&base_href, "api/v1/user/playlist/categories"),
user_playlist_bouquet_path: concat_path_leading_slash(&base_href, "api/v1/user/playlist/bouquet"),
}
}
pub async fn get_playlist_categories(&self) -> Result<Option<Rc<PlaylistCategoriesDto>>, Error> {
request_get::<Rc<PlaylistCategoriesDto>>(&self.user_playlist_categories_path, None, None)
.await
.inspect_err(|err| error!("{err}"))
}
pub async fn get_playlist_bouquet(&self) -> Result<Option<Rc<PlaylistBouquetDto>>, Error> {
request_get::<Rc<PlaylistBouquetDto>>(&self.user_playlist_bouquet_path, None, None)
.await
.inspect_err(|err| error!("{err}"))
}
pub async fn save_playlist_bouquet(&self, bouquet: &PlaylistBouquetDto) -> Result<(), Error> {
request_post::<&PlaylistBouquetDto, ()>(&self.user_playlist_bouquet_path, bouquet, None, None)
.await
.inspect_err(|err| error!("{err}"))
.map(|_| ())
}
}
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "shared"
version = "3.2.10"
version = "3.2.11"
edition = "2021"
[dependencies]
+4 -4
View File
@@ -1,4 +1,4 @@
use crate::utils::default_as_true;
use crate::utils::{default_as_true, deserialize_timestamp};
use crate::error::{TuliproxError, TuliproxErrorKind};
use crate::model::{ProxyType, ProxyUserStatus};
@@ -24,7 +24,7 @@ pub struct ProxyUserCredentialsDto {
pub epg_timeshift: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub created_at: Option<i64>,
#[serde(skip_serializing_if = "Option::is_none")]
#[serde(default, deserialize_with = "deserialize_timestamp", skip_serializing_if = "Option::is_none")]
pub exp_date: Option<i64>,
#[serde(default)]
pub max_connections: u32,
@@ -72,8 +72,8 @@ impl ProxyUserCredentialsDto {
}
}
if let Some(exp_date) = self.exp_date {
let now = chrono::Local::now();
if (exp_date - now.timestamp()) < 0 {
let now = chrono::Utc::now().timestamp();
if exp_date < now {
return false;
}
}
+7 -1
View File
@@ -1,6 +1,6 @@
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};
use crate::utils::{default_as_true, get_credentials_from_url_str, get_trimmed_string, sanitize_sensitive_info, trim_last_slash, deserialize_timestamp};
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};
@@ -241,6 +241,9 @@ pub struct ConfigInputAliasDto {
pub priority: i16,
#[serde(default)]
pub max_connections: u16,
#[serde(default, deserialize_with = "deserialize_timestamp", skip_serializing_if = "Option::is_none")]
pub exp_date: Option<i64>,
}
impl ConfigInputAliasDto {
@@ -295,6 +298,8 @@ pub struct ConfigInputDto {
pub method: InputFetchMethod,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub staged: Option<StagedInputDto>,
#[serde(default, deserialize_with = "deserialize_timestamp", skip_serializing_if = "Option::is_none")]
pub exp_date: Option<i64>,
}
impl Default for ConfigInputDto {
@@ -316,6 +321,7 @@ impl Default for ConfigInputDto {
max_connections: 0,
method: InputFetchMethod::default(),
staged: None,
exp_date: None,
}
}
}
+13
View File
@@ -37,6 +37,19 @@ impl XtreamCluster {
}
}
impl FromStr for XtreamCluster {
type Err = String;
fn from_str(s: &str) -> Result<Self, Self::Err> {
match s.to_lowercase().as_str() {
"live" => Ok(XtreamCluster::Live),
"video" | "vod" | "movie" => Ok(XtreamCluster::Video),
"series" => Ok(XtreamCluster::Series),
_ => Err(format!("Invalid XtreamCluster: {s}")),
}
}
}
impl Display for XtreamCluster {
fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
write!(f, "{}", self.as_str())
+23 -31
View File
@@ -1,31 +1,5 @@
// #[derive(Debug, serde::Serialize, serde::Deserialize, Default)]
// pub struct PlaylistCategories {
// #[serde(skip_serializing_if = "Option::is_none")]
// pub live: Option<Vec<String>>,
// #[serde(skip_serializing_if = "Option::is_none")]
// pub vod: Option<Vec<String>>,
// #[serde(skip_serializing_if = "Option::is_none")]
// pub series: Option<Vec<String>>,
// }
#[derive(Debug, serde::Serialize, serde::Deserialize, Default)]
pub struct PlaylistCategoryDto {
pub id: String,
pub name: String,
}
#[derive(Debug, serde::Serialize, serde::Deserialize, Default)]
pub struct PlaylistCategoriesDto {
#[serde(skip_serializing_if = "Option::is_none")]
pub live: Option<Vec<PlaylistCategoryDto>>,
#[serde(skip_serializing_if = "Option::is_none")]
pub vod: Option<Vec<PlaylistCategoryDto>>,
#[serde(skip_serializing_if = "Option::is_none")]
pub series: Option<Vec<PlaylistCategoryDto>>,
}
#[derive(Debug, serde::Serialize, serde::Deserialize, Default)]
pub struct TargetBouquetDto {
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, Default, PartialEq)]
pub struct PlaylistClusterCategoriesDto {
#[serde(skip_serializing_if = "Option::is_none")]
pub live: Option<Vec<String>>,
#[serde(skip_serializing_if = "Option::is_none")]
@@ -34,10 +8,28 @@ pub struct TargetBouquetDto {
pub series: Option<Vec<String>>,
}
#[derive(Debug, serde::Serialize, serde::Deserialize, Default)]
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, Default, PartialEq)]
pub struct PlaylistCategoriesDto {
#[serde(skip_serializing_if = "Option::is_none")]
pub xtream: Option<PlaylistClusterCategoriesDto>,
#[serde(skip_serializing_if = "Option::is_none")]
pub m3u: Option<PlaylistClusterCategoriesDto>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, Default, PartialEq)]
pub struct PlaylistClusterBouquetDto {
#[serde(skip_serializing_if = "Option::is_none")]
pub live: Option<Vec<String>>,
#[serde(skip_serializing_if = "Option::is_none")]
pub vod: Option<Vec<String>>,
#[serde(skip_serializing_if = "Option::is_none")]
pub series: Option<Vec<String>>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, Default, PartialEq)]
pub struct PlaylistBouquetDto {
#[serde(skip_serializing_if = "Option::is_none")]
pub xtream: Option<TargetBouquetDto>,
pub xtream: Option<PlaylistClusterBouquetDto>,
#[serde(skip_serializing_if = "Option::is_none")]
pub m3u: Option<TargetBouquetDto>,
pub m3u: Option<PlaylistClusterBouquetDto>,
}
+7 -4
View File
@@ -5,6 +5,7 @@ use crate::utils::{current_time_secs, longest};
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub struct StreamChannel {
pub target_id: u16,
pub virtual_id: u32,
pub provider_id: u32,
pub item_type: PlaylistItemType,
@@ -15,15 +16,16 @@ pub struct StreamChannel {
pub shared: bool,
}
pub fn create_stream_channel_with_type(pli: &XtreamPlaylistItem, item_type: PlaylistItemType) -> StreamChannel {
let mut stream_channel = pli.to_stream_channel();
pub fn create_stream_channel_with_type(target_id: u16, pli: &XtreamPlaylistItem, item_type: PlaylistItemType) -> StreamChannel {
let mut stream_channel = pli.to_stream_channel(target_id);
stream_channel.item_type = item_type;
stream_channel
}
impl XtreamPlaylistItem {
pub fn to_stream_channel(&self) -> StreamChannel {
pub fn to_stream_channel(&self, target_id: u16) -> StreamChannel {
StreamChannel {
target_id,
virtual_id: self.virtual_id,
provider_id: self.provider_id,
item_type: self.item_type,
@@ -37,8 +39,9 @@ impl XtreamPlaylistItem {
}
impl M3uPlaylistItem {
pub fn to_stream_channel(&self) -> StreamChannel {
pub fn to_stream_channel(&self, target_id: u16) -> StreamChannel {
StreamChannel {
target_id,
virtual_id: self.virtual_id,
provider_id: self.get_provider_id().unwrap_or_default(),
item_type: self.item_type,
+49 -11
View File
@@ -1,8 +1,9 @@
use std::io;
use serde::{Deserialize, Deserializer,};
use serde::de::DeserializeOwned;
use serde_json::Value;
use crate::error::to_io_error;
use chrono::{NaiveDateTime, ParseError, TimeZone, Utc};
use serde::de::DeserializeOwned;
use serde::Deserialize;
use serde_json::Value;
use std::io;
fn value_to_string_array(value: &[Value]) -> Vec<String> {
value.iter().filter_map(value_to_string).collect()
@@ -19,9 +20,9 @@ fn value_to_string(v: &Value) -> Option<String> {
pub fn deserialize_as_option_string<'de, D>(deserializer: D) -> Result<Option<String>, D::Error>
where
D: Deserializer<'de>,
D: serde::Deserializer<'de>,
{
let value: Value = Deserialize::deserialize(deserializer)?;
let value: Value = serde::Deserialize::deserialize(deserializer)?;
match &value {
Value::String(s) => Ok(Some(s.to_owned())),
@@ -32,9 +33,9 @@ where
pub fn deserialize_as_string<'de, D>(deserializer: D) -> Result<String, D::Error>
where
D: Deserializer<'de>,
D: serde::Deserializer<'de>,
{
let value: Value = Deserialize::deserialize(deserializer)?;
let value: Value = serde::Deserialize::deserialize(deserializer)?;
match &value {
Value::String(s) => Ok(s.to_string()),
@@ -45,7 +46,7 @@ where
pub fn deserialize_as_string_array<'de, D>(deserializer: D) -> Result<Option<Vec<String>>, D::Error>
where
D: Deserializer<'de>,
D: serde::Deserializer<'de>,
{
Value::deserialize(deserializer).map(|v| match v {
Value::String(value) => Some(vec![value]),
@@ -54,11 +55,11 @@ where
})
}
pub fn deserialize_number_from_string<'de, D, T: DeserializeOwned + std::str::FromStr>(
pub fn deserialize_number_from_string<'de, D, T: DeserializeOwned + std::str::FromStr>(
deserializer: D,
) -> Result<Option<T>, D::Error>
where
D: Deserializer<'de>,
D: serde::Deserializer<'de>,
{
// we define a local enum type inside of the function
// because it is untagged, serde will deserialize as the first variant
@@ -145,4 +146,41 @@ where
S: serde::Serializer,
{
serializer.serialize_str(&u8_16_to_hex(bytes))
}
/// Deserializes a timestamp from either a Unix timestamp (seconds) or a UTC datetime string
/// in the format "YYYY-MM-DD HH:MM:SS". Note: Datetime strings are interpreted as UTC.
pub fn deserialize_timestamp<'de, D>(deserializer: D) -> Result<Option<i64>, D::Error>
where
D: serde::Deserializer<'de>,
{
// - try to deserialize as seconds
// - try to deserialize as date-time string of format like "2028-11-23 14:12:34"
let val = Option::<Value>::deserialize(deserializer)?;
match val {
Some(Value::Number(n)) => n
.as_i64()
.ok_or_else(|| serde::de::Error::custom("invalid number"))
.map(Some),
Some(Value::String(s)) => parse_timestamp(&s).map_err(serde::de::Error::custom),
Some(Value::Null) => Ok(None),
Some(_) => Err(serde::de::Error::custom("expected number or string")),
None => Ok(None),
}
}
pub fn parse_timestamp(value: &str) -> Result<Option<i64>, ParseError> {
let value = value.trim();
if value.is_empty() {
return Ok(None);
}
if let Ok(ts) = value.parse::<i64>() {
return Ok(Some(ts));
}
// "YYYY-MM-DD HH:MM:SS"
let dt = NaiveDateTime::parse_from_str(value, "%Y-%m-%d %H:%M:%S")?;
let timestamp = Utc.from_utc_datetime(&dt).timestamp();
Ok(Some(timestamp))
}