Store FFmpeg in redis.

This commit is contained in:
SergeantPanda 2025-06-09 17:33:59 -05:00
parent 44c8189c29
commit 9d8e011e2c
2 changed files with 169 additions and 50 deletions

View file

@ -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"

View file

@ -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"""