Merge remote-tracking branch 'darkbreakpoint/warn-shared-stream-lag' into feature/shared_stream_diagnostic

This commit is contained in:
euzu
2025-11-18 21:01:10 +01:00
7 changed files with 169 additions and 37 deletions
+1
View File
@@ -6,6 +6,7 @@
- Async persist cache write pipe so response caching no longer blocks the async runtime
- M3U playlist exports now stream async to keep the runtime responsive
- Shared stream burst buffer zero copy data buffer to reduce memory usage.
- Added detailed shared-stream/buffer/provider logging to trace lag, cache persistence, and session/provider lifecycle events.
# 3.2.0 (2025-11-14)
+30 -6
View File
@@ -22,7 +22,7 @@ use axum::response::IntoResponse;
use chrono::{DateTime, Utc};
use futures::{StreamExt, TryStreamExt};
use jsonwebtoken::{decode, Algorithm, DecodingKey, Validation};
use log::{debug, error, log_enabled, trace};
use log::{debug, error, info, log_enabled, trace, warn};
use reqwest::header::RETRY_AFTER;
use serde::Serialize;
use shared::model::{Claims, InputFetchMethod, PlaylistEntry, PlaylistItemType, StreamChannel, TargetType, UserConnectionPermission, XtreamCluster};
@@ -1043,6 +1043,7 @@ async fn build_stream_response(
resource_url: &str,
response: reqwest::Response,
) -> axum::response::Response {
let sanitized_resource_url = sanitize_sensitive_info(resource_url);
let status = response.status();
let mut response_builder =
axum::response::Response::builder().status(status);
@@ -1070,18 +1071,41 @@ async fn build_stream_response(
// Cache only complete responses (200 OK without Content-Range)
let can_cache = status == axum::http::StatusCode::OK && !has_content_range;
if can_cache {
info!(
"Caching eligible resource stream {}",
sanitized_resource_url
);
let cache_resource_path = if let Some(cache) = app_state.cache.load().as_ref() {
Some(cache.lock().await.store_path(resource_url))
} else {
None
};
if let Some(resource_path) = cache_resource_path {
if let Ok(file) = create_new_file_for_write(&resource_path).await {
let writer = BufWriter::new(file);
let add_cache_content = get_add_cache_content(resource_url, &app_state.cache);
let stream = PersistPipeStream::new(byte_stream, writer, add_cache_content);
return try_unwrap_body!(response_builder.body(axum::body::Body::from_stream(stream)));
match create_new_file_for_write(&resource_path).await {
Ok(file) => {
info!(
"Persisting resource stream {} to {}",
sanitized_resource_url,
resource_path.display()
);
let writer = BufWriter::new(file);
let add_cache_content = get_add_cache_content(resource_url, &app_state.cache);
let stream = PersistPipeStream::new(byte_stream, writer, add_cache_content);
return try_unwrap_body!(response_builder.body(axum::body::Body::from_stream(stream)));
}
Err(err) => {
warn!(
"Failed to create cache file {} for {}: {err}",
resource_path.display(),
sanitized_resource_url
);
}
}
} else {
debug!(
"Resource cache unavailable; streaming response for {} without persistence",
sanitized_resource_url
);
}
}
@@ -1,7 +1,7 @@
use crate::api::model::provider_lineup_manager::{ProviderAllocation, ProviderLineupManager};
use crate::api::model::{EventManager, ProviderConfig};
use crate::model::{AppConfig, ConfigInput};
use log::{debug, error};
use log::{debug, error, info};
use shared::utils::{default_grace_period_millis, default_grace_period_timeout_secs, sanitize_sensitive_info};
use std::collections::{HashMap, HashSet};
use std::net::SocketAddr;
@@ -97,7 +97,10 @@ impl ActiveProviderManager {
old.release().await;
}
debug!("Added provider connection {provider_name:?} for {addr}");
info!(
"Added provider connection {provider_name:?} for {addr} (active_single_connections={})",
connections.single.len()
);
return Some(ProviderHandle::new(*addr, allocation));
}
}
@@ -132,7 +135,10 @@ impl ActiveProviderManager {
let handle = connections.single.remove(addr);
if let Some(allocation) = handle {
debug!("Released provider connection {:?} for {addr}", allocation.get_provider_name().unwrap_or_default());
info!(
"Released provider connection {:?} for {addr}",
allocation.get_provider_name().unwrap_or_default()
);
allocation.release().await;
}
@@ -152,6 +158,10 @@ impl ActiveProviderManager {
if released {
connections.shared.key_by_addr.remove(addr);
connections.shared.by_key.remove(&key);
info!(
"Released shared provider allocation {} after last subscriber {addr}",
sanitize_sensitive_info(&key)
);
}
}
@@ -165,6 +175,10 @@ impl ActiveProviderManager {
if let Some(allocation) = handle {
connections.shared.by_key.insert(key.to_string(), SharedAllocation { allocation, connections: HashSet::from([*addr]) });
connections.shared.key_by_addr.insert(*addr, key.to_string());
info!(
"Promoted provider connection for {addr} into shared stream {}",
sanitize_sensitive_info(key)
);
}
}
@@ -177,10 +191,15 @@ impl ActiveProviderManager {
if let Some(shared_allocation) = connections.shared.by_key.get_mut(&key) {
shared_allocation.connections.insert(*addr);
debug!(
"Added shared stream subscriber {addr} to {}, total shared connections={}",
sanitize_sensitive_info(&key),
shared_allocation.connections.len()
);
}
}
pub async fn get_provider_connections_count(&self) -> usize {
self.providers.active_connection_count().await
}
}
}
+24 -6
View File
@@ -119,7 +119,7 @@ impl ActiveUserManager {
}
pub async fn release_connection(&self, addr: &SocketAddr) {
let log_active_user = {
let (log_active_user, disconnected_user) = {
let mut user_connections = self.connections.write().await;
if let Some(username) = user_connections.key_by_addr.remove(addr) {
@@ -134,12 +134,16 @@ impl ActiveUserManager {
}
connection_data.streams.retain(|c| c.addr != *addr);
}
true
(true, Some(username))
} else {
false
(false, None)
}
};
if let Some(username) = disconnected_user {
info!("Released connection for user {username} at {addr}");
}
if log_active_user {
self.log_active_user().await;
}
@@ -304,7 +308,11 @@ impl ActiveUserManager {
user_connections
.key_by_addr
.insert(fingerprint.addr, username.to_string());
debug!("Added new connection for {username} at {}", fingerprint.addr);
info!(
"Added new connection for {username} at {} (active_user_connections={})",
fingerprint.addr,
connection_data.connections
);
stream_info
}
@@ -341,7 +349,10 @@ impl ActiveUserManager {
let username = user.username.clone();
let mut user_connections = self.connections.write().await;
let connection_data = user_connections.by_key.entry(username.clone()).or_insert_with(|| {
debug!("Creating session for user {username} with token {session_token} {}", sanitize_sensitive_info(stream_url));
info!(
"Creating first session for user {username} {}",
sanitize_sensitive_info(stream_url)
);
let mut data = UserConnectionData::new(0, user.max_connections);
let session = Self::new_user_session(session_token, virtual_id, provider, stream_url, addr, connection_permission);
data.add_session(session);
@@ -365,7 +376,11 @@ impl ActiveUserManager {
}
// If no session exists, create one
debug!("Creating session for user {} with token {session_token} for url: {}", user.username, sanitize_sensitive_info(stream_url));
info!(
"Creating session for user {} with token {session_token} for url: {}",
user.username,
sanitize_sensitive_info(stream_url)
);
let session = Self::new_user_session(session_token, virtual_id, provider, stream_url, addr, connection_permission);
let token = session.token.clone();
connection_data.add_session(session);
@@ -384,6 +399,9 @@ impl ActiveUserManager {
stream.addr = *addr;
}
}
info!(
"Updated session {token} for {username} address {previous_addr} -> {addr}"
);
}
}
}
@@ -4,6 +4,7 @@ use std::{
sync::Arc,
};
use std::cmp::{max};
use log::{debug, warn};
use tokio::sync::mpsc::{channel, Sender};
use tokio_stream::wrappers::ReceiverStream;
use crate::api::model::{BoxedProviderStream};
@@ -36,20 +37,35 @@ impl BufferedStream {
while client_close_signal.is_active() {
match stream.next().await {
Some(Ok(chunk)) => {
if tx.send(Ok(chunk)).await.is_err() {
client_close_signal.notify();
break;
}
let chunk_len = chunk.len();
if tx.send(Ok(chunk)).await.is_err() {
warn!(
"Buffered stream channel closed before delivering {} bytes to client",
chunk_len
);
client_close_signal.notify();
break;
}
}
Some(Err(err)) => {
let err_msg = err.to_string();
if tx.send(Err(err)).await.is_err() {
warn!(
"Buffered stream dropped stream error due to closed receiver: {err_msg}"
);
client_close_signal.notify();
}
break;
}
None => break,
None => {
debug!("Upstream provider completed buffered stream");
break;
}
}
}
if !client_close_signal.is_active() {
debug!("Client close signal fired; buffered stream exiting");
}
drop(tx);
}
}
@@ -1,6 +1,6 @@
use crate::api::model::StreamError;
use bytes::Bytes;
use log::error;
use log::{error, info, warn};
use std::collections::VecDeque;
use std::pin::Pin;
use std::sync::atomic::{AtomicUsize, Ordering};
@@ -47,7 +47,8 @@ where
fn poll_pending_writes(&mut self, cx: &mut Context<'_>) -> Poll<()> {
while let Some(chunk) = self.pending_writes.front() {
if self.current_offset >= chunk.len() {
let chunk_len = chunk.len();
if self.current_offset >= chunk_len {
if let Some(finished) = self.pending_writes.pop_front() {
self.size.fetch_add(finished.len(), Ordering::SeqCst);
}
@@ -66,6 +67,10 @@ where
self.current_offset += written;
}
Poll::Ready(Err(err)) => {
warn!(
"Dropping {} buffered bytes after persistence write error: {err}",
chunk_len.saturating_sub(self.current_offset)
);
error!("Error writing to resource file: {err}");
self.pending_writes.pop_front();
self.current_offset = 0;
@@ -92,6 +97,7 @@ where
if !self.completed {
self.completed = true;
let size = self.size.load(Ordering::SeqCst);
info!("Persisted {} bytes to cache resource", size);
(self.callback)(size);
}
}
@@ -14,7 +14,7 @@ use std::sync::Arc;
use crate::api::model::streams::buffered_stream::CHANNEL_SIZE;
use crate::api::model::BoxedProviderStream;
use crate::utils::trace_if_enabled;
use log::{debug, trace};
use log::{debug, info, trace, warn};
use shared::utils::sanitize_sensitive_info;
use std::pin::Pin;
use std::task::{Context, Poll};
@@ -98,6 +98,11 @@ impl BurstBuffer {
pub fn push(&mut self, packet: Arc<Bytes>) {
while self.current_bytes > self.buffer_size {
if let Some(popped) = self.buffer.pop_front() {
warn!(
"Shared-stream burst buffer full ({} bytes). Dropping {} bytes",
self.current_bytes,
popped.len()
);
self.current_bytes -= popped.len();
} else {
self.current_bytes = 0;
@@ -117,7 +122,7 @@ async fn send_burst_buffer(
for buf in start_buffer {
if cancellation_token.is_cancelled() { return; }
if let Err(err) = client_tx.send(buf.as_ref().clone()).await {
debug!("Error sending current chunk: {err}");
warn!("Error sending burst-buffer chunk to client: {err}");
return; // stop on send error
}
}
@@ -158,10 +163,19 @@ impl SharedStreamState {
let (client_tx, client_rx) = mpsc::channel(self.buf_size);
let mut broadcast_rx = self.broadcaster.subscribe();
let cancel_token = CancellationToken::new();
self.subscribers.write().await.insert(*addr, cancel_token.clone());
{
let mut subs = self.subscribers.write().await;
subs.insert(*addr, cancel_token.clone());
info!(
"Shared stream subscriber added {addr}; total subscribers={}",
subs.len()
);
}
let client_tx_clone = client_tx.clone();
let burst_buffer = self.burst_buffer.clone();
let burst_buffer_for_log = Arc::clone(&self.burst_buffer);
let yield_counter = YIELD_COUNTER;
let address = *addr;
tokio::spawn(async move {
@@ -188,15 +202,22 @@ impl SharedStreamState {
break;
}
loop_cnt += 1;
if loop_cnt >= YIELD_COUNTER {
if loop_cnt >= yield_counter {
tokio::task::yield_now().await;
loop_cnt = 0;
}
}
Err(tokio::sync::broadcast::error::RecvError::Lagged(skipped)) => {
trace!("Client lagged behind. Skipped {skipped} messages. {address}");
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 {} bytes, yield counter {yield_counter})",
buffered_bytes
);
loop_cnt += 1;
if loop_cnt >= YIELD_COUNTER {
if loop_cnt >= yield_counter {
tokio::task::yield_now().await;
loop_cnt = 0;
}
@@ -329,6 +350,11 @@ impl SharedStreamManager {
}
if let Some(shared_state) = shared_state {
let remaining = shared_state.subscribers.read().await.len();
info!(
"Unregistering shared stream {} (remaining_subscribers={remaining}, send_stop_signal={send_stop_signal})",
sanitize_sensitive_info(stream_url)
);
debug_if_enabled!("Unregistering shared stream {}", sanitize_sensitive_info(stream_url));
if let Some(provider_handle) = &shared_state.provider_guard {
@@ -345,23 +371,35 @@ 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) = {
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)
(
if send_stop_signal { tx } else { None },
is_empty,
subs.len(),
)
};
info!(
"Shared stream subscriber removed {addr}; remaining subscribers={remaining}"
);
if is_empty {
if let Some(url) = stream_url.as_ref() {
info!(
"No subscribers remain for {} after removing {addr}",
sanitize_sensitive_info(url)
);
self.unregister(url, send_stop_signal).await;
}
}
@@ -389,6 +427,10 @@ impl SharedStreamManager {
let mut shared_streams = self.shared_streams.write().await;
shared_streams.by_key.insert(stream_url.to_string(), shared_state);
shared_streams.key_by_addr.insert(*addr, stream_url.to_string());
info!(
"Registered shared stream {} for initial subscriber {addr}",
sanitize_sensitive_info(stream_url)
);
}
pub(crate) async fn register_shared_stream<S, E>(
@@ -411,6 +453,12 @@ impl SharedStreamManager {
app_state.active_provider.make_shared_connection(addr, stream_url).await;
let subscribed_stream = Self::subscribe_shared_stream(app_state, stream_url, addr).await;
shared_state.broadcast(stream_url, bytes_stream, Arc::clone(&app_state.shared_stream_manager));
info!(
"Created shared provider stream {} (channel_capacity={}, burst_buffer_min={} bytes)",
sanitize_sensitive_info(stream_url),
buf_size,
min_buffer_bytes
);
debug_if_enabled!("Created shared provider stream {}", sanitize_sensitive_info(stream_url));
subscribed_stream
}