diff --git a/CHANGELOG.md b/CHANGELOG.md index 24c1c0302..bcefaa58c 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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) diff --git a/backend/src/api/api_utils.rs b/backend/src/api/api_utils.rs index a086fb0b3..f1b29008a 100644 --- a/backend/src/api/api_utils.rs +++ b/backend/src/api/api_utils.rs @@ -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 + ); } } diff --git a/backend/src/api/model/active_provider_manager.rs b/backend/src/api/model/active_provider_manager.rs index 26e6a2912..9524be29d 100644 --- a/backend/src/api/model/active_provider_manager.rs +++ b/backend/src/api/model/active_provider_manager.rs @@ -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 } -} \ No newline at end of file +} diff --git a/backend/src/api/model/active_user_manager.rs b/backend/src/api/model/active_user_manager.rs index 203aa39de..b0f72d14b 100644 --- a/backend/src/api/model/active_user_manager.rs +++ b/backend/src/api/model/active_user_manager.rs @@ -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}" + ); } } } diff --git a/backend/src/api/model/streams/buffered_stream.rs b/backend/src/api/model/streams/buffered_stream.rs index 986d40367..414000f1a 100644 --- a/backend/src/api/model/streams/buffered_stream.rs +++ b/backend/src/api/model/streams/buffered_stream.rs @@ -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); } } diff --git a/backend/src/api/model/streams/persist_pipe_stream.rs b/backend/src/api/model/streams/persist_pipe_stream.rs index 550384350..92eaea05a 100644 --- a/backend/src/api/model/streams/persist_pipe_stream.rs +++ b/backend/src/api/model/streams/persist_pipe_stream.rs @@ -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); } } diff --git a/backend/src/api/model/streams/shared_stream_manager.rs b/backend/src/api/model/streams/shared_stream_manager.rs index 6c55132c3..a58b03cb4 100644 --- a/backend/src/api/model/streams/shared_stream_manager.rs +++ b/backend/src/api/model/streams/shared_stream_manager.rs @@ -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) { 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( @@ -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 }