diff --git a/apps/proxy/ts_proxy/constants.py b/apps/proxy/ts_proxy/constants.py index 4827b24b..91357bbe 100644 --- a/apps/proxy/ts_proxy/constants.py +++ b/apps/proxy/ts_proxy/constants.py @@ -63,6 +63,13 @@ class ChannelMetadataField: STREAM_SWITCH_TIME = "stream_switch_time" STREAM_SWITCH_REASON = "stream_switch_reason" + # FFmpeg performance metrics + FFMPEG_SPEED = "ffmpeg_speed" + FFMPEG_FPS = "ffmpeg_fps" + ACTUAL_FPS = "actual_fps" + FFMPEG_BITRATE = "ffmpeg_bitrate" + FFMPEG_STATS_UPDATED = "ffmpeg_stats_updated" + # Client metadata fields CONNECTED_AT = "connected_at" LAST_ACTIVE = "last_active" diff --git a/apps/proxy/ts_proxy/stream_manager.py b/apps/proxy/ts_proxy/stream_manager.py index 973f86b9..7ba5f9f8 100644 --- a/apps/proxy/ts_proxy/stream_manager.py +++ b/apps/proxy/ts_proxy/stream_manager.py @@ -7,6 +7,7 @@ import socket import requests import subprocess import gevent # Add this import +import re # Add this import at the top from typing import Optional, List from django.shortcuts import get_object_or_404 from apps.proxy.config import TSConfig as Config @@ -376,73 +377,184 @@ class StreamManager: logger.debug(f"Started stderr reader thread for channel {self.channel_id}") def _read_stderr(self): - """Read and log ffmpeg stderr output with stats lines combined""" + """Read and log ffmpeg stderr output with real-time stats parsing""" try: buffer = b"" - in_stats_line = False + stats_buffer = b"" # Common FFmpeg stats line prefixes for detection stats_prefixes = [b"frame=", b"size=", b"time=", b"bitrate=", b"speed="] # Read in small chunks while self.transcode_process and self.transcode_process.stderr: - chunk = self.transcode_process.stderr.read(1) - if not chunk: # EOF + try: + chunk = self.transcode_process.stderr.read(256) # Smaller chunks for real-time processing + if not chunk: + break + + buffer += chunk + + # Check for stats data in the buffer (stats usually start with "frame=") + if b"frame=" in buffer: + # Split buffer at frame= to isolate potential stats + parts = buffer.split(b"frame=") + + # Process any complete lines before the stats + if parts[0]: + line_buffer = parts[0] + while b'\n' in line_buffer or b'\r' in line_buffer: + if b'\n' in line_buffer: + line, line_buffer = line_buffer.split(b'\n', 1) + else: + line, line_buffer = line_buffer.split(b'\r', 1) + + if line.strip(): + self._log_stderr_content(line.decode('utf-8', errors='ignore')) + + # Handle stats data - combine with previous stats buffer + if len(parts) > 1: + stats_buffer = b"frame=" + parts[-1] + + # Look for common stats patterns to determine if we have a complete stats line + # Stats typically contain: frame=X fps=Y q=Z size=A time=B bitrate=C speed=D + if b"speed=" in stats_buffer: + # We likely have a complete or near-complete stats line + try: + stats_text = stats_buffer.decode('utf-8', errors='ignore') + self._parse_ffmpeg_stats(stats_text) + self._log_stderr_content(stats_text) + stats_buffer = b"" # Clear stats buffer after processing + except Exception as e: + logger.debug(f"Error parsing stats: {e}") + + # Keep any remaining parts for next iteration + buffer = b"" + else: + # No stats data, process as regular lines + while b'\n' in buffer or b'\r' in buffer: + if b'\n' in buffer: + line, buffer = buffer.split(b'\n', 1) + else: + line, buffer = buffer.split(b'\r', 1) + + if line.strip(): + self._log_stderr_content(line.decode('utf-8', errors='ignore')) + + # Prevent buffer from growing too large + if len(buffer) > 4096: + # Keep only the last 1024 bytes to preserve potential incomplete data + buffer = buffer[-1024:] + + if len(stats_buffer) > 2048: # Stats lines shouldn't be this long + stats_buffer = b"" + + except Exception as e: + logger.error(f"Error reading stderr: {e}") break - buffer += chunk - - # Check if we have a complete line - if chunk == b'\n' or chunk == b'\r': - # We have a complete line - if buffer.strip(): - line = buffer.decode('utf-8', errors='replace').strip() - # Check if this is a stats line - if any(stat_prefix in line for stat_prefix in [p.decode() for p in stats_prefixes]): - self._log_stderr_content(f"FFmpeg stats: {line}") - else: - self._log_stderr_content(line) - buffer = b"" - in_stats_line = False - continue - - # Check if this is the start of a new non-stats line after we were in a stats line - if in_stats_line: - # If we see two consecutive newlines or a non-stats-related line, - # consider the stats block complete - if chunk == b'\n' or chunk == b'\r': - in_stats_line = False - if buffer.strip(): - self._log_stderr_content(buffer.decode('utf-8', errors='replace').strip()) - buffer = b"" - # Check if this is the start of a new stats line - elif any(prefix in buffer for prefix in stats_prefixes): - # We're now in a stats line - in_stats_line = True except Exception as e: # Catch any other exceptions in the thread to prevent crashes try: - logger.error(f"Error in stderr reader thread: {e}") + logger.error(f"Error in stderr reader thread for channel {self.channel_id}: {e}") except: - # Again, if logging fails, continue silently pass def _log_stderr_content(self, content): - """Helper method to log stderr content with error handling""" + """Log stderr content from FFmpeg with appropriate log levels""" try: - # Wrap the logging call in a try-except to prevent crashes due to logging errors - logger.debug(f"Transcode stderr [{self.channel_id}]: {content}") - except OSError as e: - # If logging fails, try a simplified log message - if e.errno == 105: # No buffer space available - try: - # Try a much shorter message without the error content - logger.warning(f"Logging error (buffer full) in channel {self.channel_id}") - except: - # If even that fails, we have to silently continue - pass - except Exception: - # Ignore other logging errors to prevent thread crashes - pass + content = content.strip() + if not content: + return + + # Convert to lowercase for easier matching + content_lower = content.lower() + + # Determine log level based on content + if any(keyword in content_lower for keyword in ['error', 'failed', 'cannot', 'invalid', 'corrupt']): + logger.error(f"FFmpeg stderr: {content}") + elif any(keyword in content_lower for keyword in ['warning', 'deprecated', 'ignoring']): + logger.warning(f"FFmpeg stderr: {content}") + elif content.startswith('frame=') or 'fps=' in content or 'speed=' in content: + # Stats lines - log at debug level to avoid spam + logger.debug(f"FFmpeg stats: {content}") + elif any(keyword in content_lower for keyword in ['input', 'output', 'stream', 'video', 'audio']): + # Stream info - log at info level + logger.info(f"FFmpeg info: {content}") + else: + # Everything else at debug level + logger.debug(f"FFmpeg stderr: {content}") + + except Exception as e: + logger.error(f"Error logging stderr content: {e}") + + def _parse_ffmpeg_stats(self, stats_line): + """Parse FFmpeg stats line and extract speed, fps, and bitrate""" + try: + # Example FFmpeg stats line: + # frame= 1234 fps= 30 q=28.0 size= 2048kB time=00:00:41.33 bitrate= 406.1kbits/s speed=1.02x + + # Extract speed (e.g., "speed=1.02x") + speed_match = re.search(r'speed=\s*([0-9.]+)x?', stats_line) + ffmpeg_speed = float(speed_match.group(1)) if speed_match else None + + # Extract fps (e.g., "fps= 30") + fps_match = re.search(r'fps=\s*([0-9.]+)', stats_line) + ffmpeg_fps = float(fps_match.group(1)) if fps_match else None + + # Extract bitrate (e.g., "bitrate= 406.1kbits/s") + bitrate_match = re.search(r'bitrate=\s*([0-9.]+(?:\.[0-9]+)?)\s*([kmg]?)bits/s', stats_line, re.IGNORECASE) + ffmpeg_bitrate = None + if bitrate_match: + bitrate_value = float(bitrate_match.group(1)) + unit = bitrate_match.group(2).lower() + # Convert to kbps + if unit == 'm': + bitrate_value *= 1000 + elif unit == 'g': + bitrate_value *= 1000000 + # If no unit or 'k', it's already in kbps + ffmpeg_bitrate = bitrate_value + + # Calculate actual FPS + actual_fps = None + if ffmpeg_fps is not None and ffmpeg_speed is not None and ffmpeg_speed > 0: + actual_fps = ffmpeg_fps / ffmpeg_speed + + # Store in Redis if we have valid data + if any(x is not None for x in [ffmpeg_speed, ffmpeg_fps, actual_fps, ffmpeg_bitrate]): + self._update_ffmpeg_stats_in_redis(ffmpeg_speed, ffmpeg_fps, actual_fps, ffmpeg_bitrate) + + logger.debug(f"FFmpeg stats - Speed: {ffmpeg_speed}x, FFmpeg FPS: {ffmpeg_fps}, " + f"Actual FPS: {actual_fps:.1f if actual_fps else 'N/A'}, " + f"Bitrate: {ffmpeg_bitrate:.1f if ffmpeg_bitrate else 'N/A'} kbps") + + except Exception as e: + logger.debug(f"Error parsing FFmpeg stats: {e}") + + def _update_ffmpeg_stats_in_redis(self, speed, fps, actual_fps, bitrate): + """Update FFmpeg performance stats in Redis metadata""" + try: + if hasattr(self.buffer, 'redis_client') and self.buffer.redis_client: + metadata_key = RedisKeys.channel_metadata(self.channel_id) + update_data = { + ChannelMetadataField.FFMPEG_STATS_UPDATED: str(time.time()) + } + + if speed is not None: + update_data[ChannelMetadataField.FFMPEG_SPEED] = str(round(speed, 3)) + + if fps is not None: + update_data[ChannelMetadataField.FFMPEG_FPS] = str(round(fps, 1)) + + if actual_fps is not None: + update_data[ChannelMetadataField.ACTUAL_FPS] = str(round(actual_fps, 1)) + + if bitrate is not None: + update_data[ChannelMetadataField.FFMPEG_BITRATE] = str(round(bitrate, 1)) + + self.buffer.redis_client.hset(metadata_key, mapping=update_data) + + except Exception as e: + logger.error(f"Error updating FFmpeg stats in Redis: {e}") def _establish_http_connection(self): """Establish a direct HTTP connection to the stream"""