diff --git a/CHANGELOG.md b/CHANGELOG.md index f6713ad4b..587ab28e4 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,6 +1,11 @@ # Changelog # 3.1.4 (2025-06-xx) -- fixed custom stream loop +- share live stream refactored +- fixed active user count +- fixed hls streaming +- more logs sanitized +- added session key for session management +- added sleep timer `sleep_timer_mins` to config.yml # 3.1.3 (2025-06-06) - Fixed xtream codes series info duplicate fields problem. diff --git a/src/api/model/active_user_manager.rs b/src/api/model/active_user_manager.rs index 6daf5aa50..08822138c 100644 --- a/src/api/model/active_user_manager.rs +++ b/src/api/model/active_user_manager.rs @@ -1,5 +1,6 @@ use crate::model::Config; use crate::model::{ProxyUserCredentials, UserConnectionPermission}; +use crate::utils::request::sanitize_sensitive_info; use crate::utils::{current_time_secs, default_grace_period_millis, default_grace_period_timeout_secs}; use jsonwebtoken::get_current_timestamp; use log::{debug, info}; @@ -7,7 +8,6 @@ use std::collections::HashMap; use std::sync::atomic::{AtomicU64, Ordering}; use std::sync::Arc; use tokio::sync::RwLock; -use crate::utils::request::sanitize_sensitive_info; const USER_CON_TTL: u64 = 10_800; // 3 hours const USER_SESSION_LIMIT: usize = 50; @@ -198,12 +198,10 @@ impl ActiveUserManager { if connection_data.connections == 0 { lock.remove(username); - } else { - if connection_data.connections < connection_data.max_connections { - // Grace timeout expired, reset grace counters - connection_data.granted_grace = false; - connection_data.grace_ts = 0; - } + } else if connection_data.connections < connection_data.max_connections { + // Grace timeout expired, reset grace counters + connection_data.granted_grace = false; + connection_data.grace_ts = 0; } } drop(lock); diff --git a/src/api/model/streams/active_client_stream.rs b/src/api/model/streams/active_client_stream.rs index e68147d57..5b62d36fc 100644 --- a/src/api/model/streams/active_client_stream.rs +++ b/src/api/model/streams/active_client_stream.rs @@ -14,6 +14,8 @@ use std::pin::Pin; use std::sync::atomic::AtomicU8; use std::sync::{Arc, Mutex}; use std::task::{Poll, Waker}; +use crate::api::model::streams::timed_client_stream::TimedClientStream; +use futures::{StreamExt}; const INNER_STREAM: u8 = 0_u8; const GRACE_BLOCK_STREAM: u8 = 1_u8; @@ -55,8 +57,21 @@ impl ActiveClientStream { c.provider_connections_exhausted.clone() )); + let stream = stream_details.stream.take().unwrap(); + let stream = match app_state.config.sleep_timer_mins { + None => stream, + Some(mins) => { + let secs = u32::try_from((u64::from(mins) * 60).min(u64::from(u32::MAX))).unwrap_or(0); + if secs > 0 { + TimedClientStream::new(stream, secs).boxed() + } else { + stream + } + } + }; + Self { - inner: stream_details.stream.take().unwrap(), + inner: stream, user_connection_guard, provider_connection_guard: stream_details.provider_connection_guard, send_custom_stream_flag: grace_stop_flag, diff --git a/src/api/model/streams/provider_stream_factory.rs b/src/api/model/streams/provider_stream_factory.rs index e1a6e5954..250d19f0c 100644 --- a/src/api/model/streams/provider_stream_factory.rs +++ b/src/api/model/streams/provider_stream_factory.rs @@ -5,7 +5,7 @@ use crate::api::model::stream_error::StreamError; use crate::api::model::streams::buffered_stream::BufferedStream; use crate::api::model::streams::client_stream::ClientStream; use crate::api::model::streams::provider_stream::{create_channel_unavailable_stream, get_header_filter_for_item_type}; -use crate::api::model::streams::timed_client_stream::TimeoutClientStream; +use crate::api::model::streams::timed_client_stream::TimedClientStream; use crate::model::PlaylistItemType; use crate::model::{Config, DEFAULT_USER_AGENT}; use crate::tools::atomic_once_flag::AtomicOnceFlag; @@ -260,7 +260,7 @@ async fn provider_stream_request(cfg: &Config, request_client: Arc 0 { - TimeoutClientStream::new(provider_stream, stream_options.get_reconnect_force_secs()).boxed() + TimedClientStream::new(provider_stream, stream_options.get_reconnect_force_secs()).boxed() } else { provider_stream }; diff --git a/src/api/model/streams/shared_stream_manager.rs b/src/api/model/streams/shared_stream_manager.rs index 58be8c60e..f6f4f2166 100644 --- a/src/api/model/streams/shared_stream_manager.rs +++ b/src/api/model/streams/shared_stream_manager.rs @@ -1,25 +1,25 @@ use crate::api::model::app_state::AppState; -use crate::api::model::streams::provider_stream_factory::STREAM_QUEUE_SIZE; use crate::api::model::stream_error::StreamError; +use crate::api::model::streams::provider_stream_factory::STREAM_QUEUE_SIZE; use crate::utils::debug_if_enabled; use crate::utils::request::sanitize_sensitive_info; use bytes::Bytes; use futures::stream::{BoxStream, FuturesUnordered}; use futures::{Stream, StreamExt}; use std::collections::HashMap; -use std::sync::{Arc}; +use std::sync::Arc; +use tokio::sync::mpsc::Sender; use tokio::sync::RwLock; -use tokio::sync::mpsc::{Sender}; +use crate::api::model::stream::BoxedProviderStream; +use dashmap::DashMap; +use log::trace; use std::pin::Pin; use std::sync::atomic::{AtomicUsize, Ordering}; use std::task::{Context, Poll}; -use dashmap::DashMap; -use log::{trace}; -use tokio::sync::{mpsc}; use tokio::sync::mpsc::error::TrySendError; +use tokio::sync::mpsc; use tokio_stream::wrappers::ReceiverStream; -use crate::api::model::stream::BoxedProviderStream; /// /// Wraps a `ReceiverStream` as Stream> @@ -84,7 +84,7 @@ impl SharedStreamState { shared_streams: Arc, ) where - S: Stream> + Unpin + 'static + Send, + S: Stream> + Unpin + 'static + Send, E: std::fmt::Debug + Send, { let mut source_stream = Box::pin(bytes_stream); @@ -111,7 +111,7 @@ impl SharedStreamState { Ok(()) => true, Err(TrySendError::Closed(_)) => false, Err(err) => { - trace!("broadcast try_send error: {:?}", err); + trace!("broadcast try_send error: {err:?}"); true } }); @@ -135,14 +135,9 @@ impl SharedStreamState { let mut has_fillable_subscriber = false; while let Some(result) = futures.next().await { - match result { - Ok(_) => { - has_fillable_subscriber = true; - break; - } - Err(_) => { - // ignore; continue waiting - } + if let Ok(()) = result { + has_fillable_subscriber = true; + break; } } @@ -157,7 +152,7 @@ impl SharedStreamState { Ok(()) => true, Err(TrySendError::Closed(_)) => false, Err(err) => { - trace!("broadcast try_send error after reserve: {:?}", err); + trace!("broadcast try_send error after reserve: {err:?}"); true } }); @@ -196,7 +191,7 @@ impl SharedStreamManager { } async fn register(&self, stream_url: &str, shared_state: SharedStreamState) { - let _= self.shared_streams.write().await.insert(stream_url.to_string(), shared_state); + let _ = self.shared_streams.write().await.insert(stream_url.to_string(), shared_state); } pub(crate) async fn subscribe( @@ -204,10 +199,10 @@ impl SharedStreamManager { stream_url: &str, bytes_stream: S, headers: Vec<(String, String)>, - buffer_size: usize,) -> Option + buffer_size: usize, ) -> Option where S: Stream> + Unpin + 'static + std::marker::Send, - E: std::fmt::Debug + std::marker::Send + E: std::fmt::Debug + std::marker::Send, { let buf_size = std::cmp::max(buffer_size, STREAM_QUEUE_SIZE); let shared_state = SharedStreamState::new(headers, buf_size); diff --git a/src/api/model/streams/timed_client_stream.rs b/src/api/model/streams/timed_client_stream.rs index 293f090aa..f76d8b84a 100644 --- a/src/api/model/streams/timed_client_stream.rs +++ b/src/api/model/streams/timed_client_stream.rs @@ -6,22 +6,22 @@ use std::task::Poll; use std::time::{Duration, Instant}; use crate::api::model::stream::BoxedProviderStream; -pub struct TimeoutClientStream { +pub struct TimedClientStream { inner: BoxedProviderStream, - duration: Duration, - start_time: Instant, + deadline: Instant, } -impl TimeoutClientStream { +impl TimedClientStream { pub(crate) fn new(inner: BoxedProviderStream, duration: u32) -> Self { - Self { inner, duration: Duration::from_secs(u64::from(duration)) , start_time: Instant::now() } + let deadline = Instant::now() + Duration::from_secs(u64::from(duration)); + Self { inner, deadline } } } -impl Stream for TimeoutClientStream { +impl Stream for TimedClientStream { type Item = Result; fn poll_next(mut self: Pin<&mut Self>,cx: &mut std::task::Context<'_>,) -> Poll> { - if self.start_time.elapsed() > self.duration { + if Instant::now() >= self.deadline { return Poll::Ready(None); } Pin::as_mut(&mut self.inner).poll_next(cx) diff --git a/src/model/config.rs b/src/model/config.rs index 5c36138c4..26a2a17a0 100644 --- a/src/model/config.rs +++ b/src/model/config.rs @@ -169,6 +169,8 @@ pub struct ConfigDto { pub user_access_control: bool, #[serde(default = "default_connect_timeout_secs")] pub connect_timeout_secs: u32, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub sleep_timer_mins: Option, #[serde(default)] pub config_hot_reload: bool, #[serde(default, skip_serializing_if = "Option::is_none")] @@ -378,6 +380,8 @@ pub struct Config { pub user_access_control: bool, #[serde(default = "default_connect_timeout_secs")] pub connect_timeout_secs: u32, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub sleep_timer_mins: Option, #[serde(default)] pub update_on_boot: bool, #[serde(default)]