From a4c69710b713b2a99801714c2bbde7ed247341c9 Mon Sep 17 00:00:00 2001 From: SergeantPanda Date: Sat, 15 Mar 2025 17:23:06 -0500 Subject: [PATCH] Stream switching is working for both proxy and transcode --- apps/proxy/config.py | 2 +- apps/proxy/ts_proxy/server.py | 4 +- apps/proxy/ts_proxy/stream_manager.py | 57 ++++++++++++++++++++++----- apps/proxy/ts_proxy/views.py | 18 +++++++-- 4 files changed, 64 insertions(+), 17 deletions(-) diff --git a/apps/proxy/config.py b/apps/proxy/config.py index d8fabfd6..ffa0b559 100644 --- a/apps/proxy/config.py +++ b/apps/proxy/config.py @@ -29,7 +29,7 @@ class TSConfig(BaseConfig): MAX_RETRIES = 3 # maximum connection retry attempts # Buffer settings - INITIAL_BEHIND_CHUNKS = 3 # How many chunks behind to start a client + INITIAL_BEHIND_CHUNKS = 4 # How many chunks behind to start a client (4 chunks = ~1MB) CHUNK_BATCH_SIZE = 5 # How many chunks to fetch in one batch KEEPALIVE_INTERVAL = 0.5 # Seconds between keepalive packets when at buffer head diff --git a/apps/proxy/ts_proxy/server.py b/apps/proxy/ts_proxy/server.py index d450f7bb..c55e7f9a 100644 --- a/apps/proxy/ts_proxy/server.py +++ b/apps/proxy/ts_proxy/server.py @@ -281,7 +281,7 @@ class ProxyServer: logger.error(f"Error extending ownership: {e}") return False - def initialize_channel(self, url, channel_id, user_agent=None, transcode_cmd=[]): + def initialize_channel(self, url, channel_id, user_agent=None, transcode=False): """Initialize a channel without redundant active key""" try: # Create buffer and client manager instances @@ -381,7 +381,7 @@ class ProxyServer: self.stream_buffers[channel_id] = buffer # Only the owner worker creates the actual stream manager - stream_manager = StreamManager(channel_url, buffer, user_agent=channel_user_agent, transcode_cmd=transcode_cmd) + stream_manager = StreamManager(channel_url, buffer, user_agent=channel_user_agent, transcode=False) logger.debug(f"Created StreamManager for channel {channel_id}") self.stream_managers[channel_id] = stream_manager diff --git a/apps/proxy/ts_proxy/stream_manager.py b/apps/proxy/ts_proxy/stream_manager.py index 44df13c2..2ba912e7 100644 --- a/apps/proxy/ts_proxy/stream_manager.py +++ b/apps/proxy/ts_proxy/stream_manager.py @@ -6,9 +6,11 @@ import time import requests import subprocess from typing import Optional, List +from django.shortcuts import get_object_or_404 from apps.proxy.config import TSConfig as Config - -# Import StreamBuffer but use lazy imports inside methods if needed +from apps.channels.models import Channel, Stream +from apps.m3u.models import M3UAccount, M3UAccountProfile +from core.models import UserAgent, CoreSettings from .stream_buffer import StreamBuffer logger = logging.getLogger("ts_proxy") @@ -16,7 +18,7 @@ logger = logging.getLogger("ts_proxy") class StreamManager: """Manages a connection to a TS stream without using raw sockets""" - def __init__(self, url, buffer, user_agent=None, transcode_cmd=[]): + def __init__(self, url, buffer, user_agent=None, transcode=False): # Basic properties self.url = url self.buffer = buffer @@ -26,10 +28,11 @@ class StreamManager: self.max_retries = Config.MAX_RETRIES self.current_response = None self.current_session = None + self.url_switching = False # Sockets used for transcode jobs self.socket = None - self.transcode_cmd = transcode_cmd + self.transcode = transcode self.transcode_process = None # User agent for connection @@ -80,7 +83,16 @@ class StreamManager: logger.info(f"Starting stream for URL: {self.url}") while self.running: - if len(self.transcode_cmd) > 0: + if self.transcode: + if self.url_switching: + logger.debug("Skipping connection attempt during URL switch") + time.sleep(.1) + continue + # Generate transcode command + logger.debug(f"Building transcode command for channel {self.channel_id}") + channel = get_object_or_404(Channel, pk=self.channel_id) + stream_profile = channel.get_stream_profile() + self.transcode_cmd = stream_profile.build_command(self.url, self.user_agent) # Start command process for transcoding logger.debug(f"Starting transcode process: {self.transcode_cmd}") self.transcode_process = subprocess.Popen( @@ -107,6 +119,10 @@ class StreamManager: else: try: # Using direct HTTP streaming + if self.url_switching: + logger.debug("Skipping connection attempt during URL switch") + time.sleep(.1) + continue logger.debug(f"Using TS Proxy to connect to stream: {self.url}") # Create new session for each connection attempt session = self._create_session() @@ -151,12 +167,13 @@ class StreamManager: self.buffer.redis_client.set(last_data_key, str(time.time()), ex=60) except (AttributeError, ConnectionError) as e: if self.stop_requested: - # This is expected during shutdown, just log at debug level logger.debug(f"Expected connection error during shutdown: {e}") + elif hasattr(self, 'url_switching') and self.url_switching: + # This is expected during URL switching, just log at debug level + logger.debug(f"Expected connection error during URL switch: {e}") else: # Unexpected error during normal operation logger.error(f"Unexpected stream error: {e}") - # Handle the error appropriately except Exception as e: # Handle the specific 'NoneType' object has no attribute 'read' error if "'NoneType' object has no attribute 'read'" in str(e): @@ -254,15 +271,23 @@ class StreamManager: self.running = False def update_url(self, new_url): - """Update stream URL and reconnect with HTTP streaming approach""" + """Update stream URL and reconnect with proper cleanup for both HTTP and transcode sessions""" if new_url == self.url: logger.info(f"URL unchanged: {new_url}") return False logger.info(f"Switching stream URL from {self.url} to {new_url}") - # Close existing HTTP connection resources instead of socket - self._close_connection() # Use our new method instead of _close_socket + # CRITICAL: Set a flag to prevent immediate reconnection with old URL + self.url_switching = True + + # Check which type of connection we're using and close it properly + if self.transcode or self.socket: + logger.debug("Closing transcode process before URL change") + self._close_socket() + else: + logger.debug("Closing HTTP connection before URL change") + self._close_connection() # Update URL and reset connection state old_url = self.url @@ -272,6 +297,18 @@ class StreamManager: # Reset retry counter to allow immediate reconnect self.retry_count = 0 + # Also reset buffer position to prevent stale data after URL change + if hasattr(self.buffer, 'reset_buffer_position'): + try: + self.buffer.reset_buffer_position() + logger.debug("Reset buffer position for clean URL switch") + except Exception as e: + logger.warning(f"Failed to reset buffer position: {e}") + + # Done with URL switch + self.url_switching = False + logger.info(f"Stream switch completed for channel {self.buffer.channel_id}") + return True def should_retry(self) -> bool: diff --git a/apps/proxy/ts_proxy/views.py b/apps/proxy/ts_proxy/views.py index 922f462a..755c3229 100644 --- a/apps/proxy/ts_proxy/views.py +++ b/apps/proxy/ts_proxy/views.py @@ -78,11 +78,18 @@ def stream_ts(request, channel_id): stream_profile = channel.get_stream_profile() if stream_profile.is_redirect(): return HttpResponseRedirect(stream_url) - - transcode_cmd = stream_profile.build_command(stream_url, stream_user_agent or "") + # Need to check if profile is transcoded + logger.debug(f"Using profile {stream_profile} for stream {stream_id}") + if stream_profile == 'PROXY' or stream_profile is None: + transcode = True + else: + transcode = True # Initialize channel with the stream's user agent (not the client's) - success = proxy_server.initialize_channel(stream_url, channel_id, stream_user_agent, transcode_cmd) + success = proxy_server.initialize_channel(stream_url, channel_id, stream_user_agent, transcode) + if proxy_server.redis_client: + metadata_key = f"ts_proxy:channel:{channel_id}:metadata" + proxy_server.redis_client.hset(metadata_key, "transcode", "1" if transcode else "0") if not success: return JsonResponse({'error': 'Failed to initialize channel'}, status=500) @@ -117,14 +124,17 @@ def stream_ts(request, channel_id): metadata_key = f"ts_proxy:channel:{channel_id}:metadata" url_bytes = proxy_server.redis_client.hget(metadata_key, "url") ua_bytes = proxy_server.redis_client.hget(metadata_key, "user_agent") + transcode_bytes = proxy_server.redis_client.hget(metadata_key, "transcode") if url_bytes: url = url_bytes.decode('utf-8') if ua_bytes: stream_user_agent = ua_bytes.decode('utf-8') + # Extract transcode setting from Redis + use_transcode = transcode_bytes == b'1' if transcode_bytes else False # Use client_user_agent as fallback if stream_user_agent is None - success = proxy_server.initialize_channel(url, channel_id, stream_user_agent or client_user_agent) + success = proxy_server.initialize_channel(url, channel_id, stream_user_agent or client_user_agent, use_transcode) if not success: logger.error(f"[{client_id}] Failed to initialize channel {channel_id} locally") return JsonResponse({'error': 'Failed to initialize channel locally'}, status=500)