From 0607957f67179c639a868f47cdf5dafeaab7e37c Mon Sep 17 00:00:00 2001 From: dekzter Date: Thu, 13 Mar 2025 19:22:35 -0400 Subject: [PATCH] integrating proxy with rest of the application, added transcoding to ts proxy --- apps/channels/models.py | 73 +++- apps/proxy/ts_proxy/server.py | 646 ++++++++++++++++++---------------- apps/proxy/ts_proxy/urls.py | 4 +- apps/proxy/ts_proxy/views.py | 225 ++++++------ core/models.py | 8 + 5 files changed, 541 insertions(+), 415 deletions(-) diff --git a/apps/channels/models.py b/apps/channels/models.py index bd4045ba..5507b1be 100644 --- a/apps/channels/models.py +++ b/apps/channels/models.py @@ -1,6 +1,8 @@ from django.db import models from django.core.exceptions import ValidationError -from core.models import StreamProfile +from django.conf import settings +from core.models import StreamProfile, CoreSettings +from core.utils import redis_client # If you have an M3UAccount model in apps.m3u, you can still import it: from apps.m3u.models import M3UAccount @@ -100,6 +102,75 @@ class Channel(models.Model): def __str__(self): return f"{self.channel_number} - {self.channel_name}" + def get_stream_profile(self): + stream_profile = self.stream_profile + if not stream_profile: + stream_profile = StreamProfile.objects.get(id=CoreSettings.objects.get(key="default-stream-profile").value) + + return stream_profile + + def get_stream(self): + """ + Finds an available stream for the requested channel and returns the selected stream and profile. + """ + + # 2. Check if a stream is already active for this channel + stream_id = redis_client.get(f"channel_stream:{self.id}") + if stream_id: + profile_id = redis_client.get(f"stream_profile:{stream_id}") + if profile_id: + return stream_id, profile_id + + # 3. Iterate through channel streams and their profiles + for stream in self.streams.all().order_by('channelstream__order'): + # Retrieve the M3U account associated with the stream. + m3u_account = stream.m3u_account + m3u_profiles = m3u_account.profiles.all() + default_profile = next((obj for obj in m3u_profiles if obj.is_default), None) + profiles = [default_profile] + [obj for obj in m3u_profiles if not obj.is_default] + + for profile in profiles: + # Skip inactive profiles + if profile.is_active == False: + continue + + profile_connections_key = f"profile_connections:{profile.id}" + current_connections = int(redis_client.get(profile_connections_key) or 0) + + # Check if profile has available slots (or unlimited connections) + if profile.max_streams == 0 or current_connections < profile.max_streams: + # Start a new stream + redis_client.set(f"channel_stream:{self.id}", stream.url) + redis_client.set(f"stream_profile:{stream.url}", profile.id) # Store only the matched profile + + # Increment connection count for profiles with limits + if profile.max_streams > 0: + redis_client.incr(profile_connections_key) + + return stream.id, profile.id # Return newly assigned stream and matched profile + + # 4. No available streams + return None, None + + def release_stream(self): + """ + Called when a stream is finished to release the lock. + """ + stream_id = redis_client.get(f"channel_stream:{self.id}") + if stream_id: + redis_client.delete(f"channel_stream:{self.id}") # Remove active stream + + # Get the matched profile for cleanup + profile_id = redis_client.get(f"stream_profile:{stream_id}") + if profile_id: + profile_connections_key = f"profile_connections:{profile_id}" + + # Only decrement if the profile had a max_connections limit + current_count = int(redis_client.get(profile_connections_key) or 0) + if current_count > 0: + redis_client.decr(profile_connections_key) + + redis_client.delete(f"stream_profile:{stream_id}") # Remove profile association class ChannelGroup(models.Model): name = models.CharField(max_length=100, unique=True) diff --git a/apps/proxy/ts_proxy/server.py b/apps/proxy/ts_proxy/server.py index 12ceb045..5f0f22f2 100644 --- a/apps/proxy/ts_proxy/server.py +++ b/apps/proxy/ts_proxy/server.py @@ -10,14 +10,16 @@ Handles live TS stream proxying with support for: import requests import threading import logging -import socket +import socket import random +import subprocess from collections import deque import time import sys from typing import Optional, Set, Deque, Dict import json from apps.proxy.config import TSConfig as Config +from apps.channels.models import Channel # Configure root logger for this module logging.basicConfig( @@ -31,32 +33,35 @@ print("TS PROXY SERVER MODULE LOADED", file=sys.stderr) class StreamManager: """Manages a connection to a TS stream with continuity tracking""" - - def __init__(self, url, buffer, user_agent=None): + + def __init__(self, url, buffer, user_agent=None, transcode_cmd=[]): # Existing initialization code self.url = url self.buffer = buffer self.running = True self.connected = False self.socket = None + self.transcode_process = None self.ready_event = threading.Event() self.retry_count = 0 self.max_retries = Config.MAX_RETRIES - + # User agent for connection self.user_agent = user_agent or Config.DEFAULT_USER_AGENT - + + self.transcode_cmd = transcode_cmd + # TS packet handling self.TS_PACKET_SIZE = 188 self.recv_buffer = bytearray() self.sync_found = False self.continuity_counters = {} - + # Stream health monitoring self.last_data_time = time.time() self.healthy = True self.health_check_interval = Config.HEALTH_CHECK_INTERVAL - + # Buffer management self._last_buffer_check = time.time() logging.info(f"Initialized stream manager for channel {buffer.channel_id}") @@ -74,15 +79,15 @@ class StreamManager: """Update stream URL and reconnect""" if new_url == self.url: return False - + logging.info(f"Switching stream URL from {self.url} to {new_url}") self.url = new_url self.connected = False self._close_socket() # Close existing connection - + # Signal health monitor to reconnect immediately self.last_data_time = 0 - + return True def should_retry(self) -> bool: @@ -99,93 +104,93 @@ class StreamManager: """Process TS packets with improved resync capability""" try: # Enhanced sync byte detection with re-sync capability - if (not self.sync_found or + if (not self.sync_found or (len(self.recv_buffer) >= 188 and self.recv_buffer[0] != 0x47)): - + # Need to find sync pattern if we haven't found it yet or lost sync if len(self.recv_buffer) >= 376: # Need at least 2 packet lengths sync_found = False - + # Look for at least two sync bytes (0x47) at 188-byte intervals for i in range(min(188, len(self.recv_buffer) - 188)): - if (self.recv_buffer[i] == 0x47 and + if (self.recv_buffer[i] == 0x47 and self.recv_buffer[i + 188] == 0x47): - + # If already had sync but lost it, log the issue if self.sync_found: logging.warning(f"Re-syncing TS stream at position {i} (lost sync)") else: logging.debug(f"TS sync found at position {i}") - + # Trim buffer to start at first sync byte self.recv_buffer = self.recv_buffer[i:] self.sync_found = True sync_found = True break - + # If we couldn't find sync in this buffer, discard partial data if not sync_found: logging.warning(f"Failed to find sync pattern - discarding {len(self.recv_buffer) - 188} bytes") if len(self.recv_buffer) > 188: self.recv_buffer = self.recv_buffer[-188:] # Keep last chunk for next attempt return False - + # If we don't have a complete packet yet, wait for more data if len(self.recv_buffer) < 188: return False - + # Calculate how many complete packets we have packet_count = len(self.recv_buffer) // 188 - + if packet_count == 0: return False - + # Verify all packets have sync bytes all_synced = True for i in range(0, packet_count): if self.recv_buffer[i * 188] != 0x47: all_synced = False break - + # If not all packets are synced, re-scan for sync if not all_synced: self.sync_found = False # Force re-sync on next call return False - + # Extract complete packets packets = self.recv_buffer[:packet_count * 188] - + # Keep remaining data in buffer self.recv_buffer = self.recv_buffer[packet_count * 188:] - + # Send packets to buffer if packets: # Log first and last sync byte to validate alignment first_sync = packets[0] if len(packets) > 0 else None last_sync = packets[188 * (packet_count - 1)] if packet_count > 0 else None - + if first_sync != 0x47 or last_sync != 0x47: logging.warning(f"TS packet alignment issue: first_sync=0x{first_sync:02x}, last_sync=0x{last_sync:02x}") # Don't process misaligned packets return False - + before_index = self.buffer.index success = self.buffer.add_chunk(bytes(packets)) after_index = self.buffer.index - + # Log successful write - if not success: + if not success: logging.warning("Failed to add chunk to buffer") - + # If successful, update last data timestamp in Redis if success and hasattr(self.buffer, 'redis_client') and self.buffer.redis_client: last_data_key = f"ts_proxy:channel:{self.buffer.channel_id}:last_data" self.buffer.redis_client.set(last_data_key, str(time.time()), ex=60) # 1 minute expiry - + return success - + return False - + except Exception as e: logging.error(f"Error processing TS packets: {e}", exc_info=True) self.sync_found = False # Reset sync state on error @@ -195,10 +200,10 @@ class StreamManager: """Process received data and add to buffer""" if not chunk: return False - + # Add to existing buffer self.recv_buffer.extend(chunk) - + # Process complete packets now return self._process_complete_packets() @@ -219,50 +224,38 @@ class StreamManager: self.connected = True self.healthy = True return - + # Start health monitor thread health_thread = threading.Thread(target=self._monitor_health, daemon=True) health_thread.start() - + current_response = None # Track the current response object current_session = None # Track the current session - + # Establish network connection import socket import requests - + logging.info(f"Starting stream for URL: {self.url}") - + while self.running: try: # Parse URL - if self.url.startswith("http"): + if self.url.startswith("http") and len(self.transcode_cmd) == 0: # HTTP connection session = self._create_session() current_session = session - + try: # Create an initial connection to get socket response = session.get(self.url, stream=True) current_response = response - + if response.status_code == 200: self.connected = True self.socket = response.raw._fp.fp.raw self.healthy = True logging.info("Successfully connected to stream source") - - # Connection successful - START GRACE PERIOD HERE - self._set_waiting_for_clients() - - # Main fetch loop - while self.running and self.connected: - if self.fetch_chunk(): - self.last_data_time = time.time() - else: - if not self.running: - break - time.sleep(0.1) else: logging.error(f"Failed to connect to stream: HTTP {response.status_code}") time.sleep(2) @@ -275,32 +268,56 @@ class StreamManager: except Exception as e: logging.debug(f"Error closing response: {e}") current_response = None - + if current_session: try: current_session.close() except Exception as e: logging.debug(f"Error closing session: {e}") current_session = None + elif len(self.transcode_cmd) > 0: + # Transcode + + self.transcode_process = subprocess.Popen( + self.transcode_cmd, + stdout=subprocess.PIPE, + stderr=subprocess.DEVNULL, # Suppress FFmpeg logs + bufsize=188 * 64 # Buffer optimized for TS packets + ) + self.socket = self.transcode_process.stdout # Read from FFmpeg output + self.connected = True else: logging.error(f"Unsupported URL scheme: {self.url}") - + + if self.socket is not None: + # Connection successful - START GRACE PERIOD HERE + self._set_waiting_for_clients() + + # Main fetch loop + while self.running and self.connected: + if self.fetch_chunk(): + self.last_data_time = time.time() + else: + if not self.running: + break + time.sleep(0.1) + # Connection retry logic if self.running and not self.connected: self.retry_count += 1 if self.retry_count > self.max_retries: logging.error(f"Maximum retry attempts ({self.max_retries}) exceeded") break - + timeout = min(2 ** self.retry_count, 30) logging.info(f"Reconnecting in {timeout} seconds... (attempt {self.retry_count})") time.sleep(timeout) - + except Exception as e: logging.error(f"Connection error: {e}") self._close_socket() time.sleep(5) - + except Exception as e: logging.error(f"Stream error: {e}") self._close_socket() @@ -308,7 +325,7 @@ class StreamManager: # Final cleanup self._close_socket() logging.info("Stream manager stopped") - + def _monitor_health(self): """Monitor stream health and attempt recovery if needed""" while self.running: @@ -319,7 +336,7 @@ class StreamManager: if self.healthy: logging.warning("Stream health check: No data received for 10+ seconds") self.healthy = False - + # After 30 seconds with no data, force reconnection if now - self.last_data_time > 30: logging.warning("Stream appears dead, forcing reconnection") @@ -330,12 +347,12 @@ class StreamManager: # Stream is receiving data again after being unhealthy logging.info("Stream health restored, receiving data again") self.healthy = True - + except Exception as e: logging.error(f"Error in health monitor: {e}") - + time.sleep(self.health_check_interval) - + def _close_socket(self): """Close the socket connection safely""" if self.socket: @@ -346,12 +363,16 @@ class StreamManager: pass self.socket = None self.connected = False + if self.transcode_process: + self.transcode_process.terminate() + self.transcode_process.wait() + self.transcode_process = None def fetch_chunk(self): """Fetch data from socket with improved buffer management""" if not self.connected or not self.socket: return False - + try: # SocketIO objects use read instead of recv and don't support settimeout try: @@ -360,21 +381,21 @@ class StreamManager: chunk = self.socket.recv(188 * 64) # Standard socket else: chunk = self.socket.read(188 * 64) # SocketIO object - + except AttributeError: # Fall back to read() if recv() isn't available chunk = self.socket.read(188 * 64) - + if not chunk: # Connection closed by server logging.warning("Server closed connection") self._close_socket() self.connected = False return False - + # Process this chunk self._process_ts_data(chunk) - + # Memory management - clear any internal buffers periodically current_time = time.time() if current_time - self._last_buffer_check > 60: # Check every minute @@ -384,16 +405,16 @@ class StreamManager: # Keep only recent data, aligned to TS packet boundary keep_size = 188 * 128 # Keep reasonable buffer self.recv_buffer = self.recv_buffer[-keep_size:] - + return True - + except (socket.timeout, socket.error) as e: # Socket error logging.error(f"Socket error: {e}") self._close_socket() self.connected = False return False - + except Exception as e: logging.error(f"Error in fetch_chunk: {e}") return False @@ -404,13 +425,13 @@ class StreamManager: if hasattr(self.buffer, 'channel_id') and hasattr(self.buffer, 'redis_client'): channel_id = self.buffer.channel_id redis_client = self.buffer.redis_client - + if channel_id and redis_client: current_time = str(time.time()) - + # SIMPLIFIED: Always use direct Redis update for reliability metadata_key = f"ts_proxy:channel:{channel_id}:metadata" - + # Check current state first current_state = None try: @@ -419,7 +440,7 @@ class StreamManager: current_state = metadata[b'state'].decode('utf-8') except Exception as e: logging.error(f"Error checking current state: {e}") - + # Only update if not already past connecting if not current_state or current_state in ["initializing", "connecting"]: # Update directly - don't rely on proxy_server reference @@ -429,7 +450,7 @@ class StreamManager: "state_changed_at": current_time } redis_client.hset(metadata_key, mapping=update_data) - + # Get configured grace period or default grace_period = getattr(Config, 'CHANNEL_INIT_GRACE_PERIOD', 20) logging.info(f"STREAM MANAGER: Updated channel {channel_id} state: {current_state or 'None'} → waiting_for_clients") @@ -441,20 +462,20 @@ class StreamManager: class StreamBuffer: """Manages stream data buffering using Redis for persistence""" - + def __init__(self, channel_id=None, redis_client=None): self.channel_id = channel_id self.redis_client = redis_client self.lock = threading.Lock() self.index = 0 self.TS_PACKET_SIZE = 188 - + # STANDARDIZED KEYS: Move buffer keys under channel namespace self.buffer_index_key = f"ts_proxy:channel:{channel_id}:buffer:index" self.buffer_prefix = f"ts_proxy:channel:{channel_id}:buffer:chunk:" - + self.chunk_ttl = getattr(Config, 'REDIS_CHUNK_TTL', 60) - + # Initialize from Redis if available if self.redis_client and channel_id: try: @@ -464,12 +485,12 @@ class StreamBuffer: logging.info(f"Initialized buffer from Redis with index {self.index}") except Exception as e: logging.error(f"Error initializing buffer from Redis: {e}") - + def add_chunk(self, chunk): """Add a chunk to the buffer""" if not chunk: return False - + try: # Ensure chunk is properly aligned with TS packets if len(chunk) % self.TS_PACKET_SIZE != 0: @@ -478,7 +499,7 @@ class StreamBuffer: if (aligned_size == 0): return False chunk = chunk[:aligned_size] - + with self.lock: # Increment index atomically if self.redis_client: @@ -486,7 +507,7 @@ class StreamBuffer: chunk_index = self.redis_client.incr(self.buffer_index_key) chunk_key = f"{self.buffer_prefix}{chunk_index}" self.redis_client.setex(chunk_key, self.chunk_ttl, chunk) - + # Update local tracking of position only self.index = chunk_index return True @@ -494,39 +515,39 @@ class StreamBuffer: # No Redis - can't function in multi-worker mode logging.error("Redis not available, cannot store chunks") return False - + except Exception as e: logging.error(f"Error adding chunk to buffer: {e}") return False - + def get_chunks(self, start_index=None): """Get chunks from the buffer with detailed logging""" try: request_id = f"req_{random.randint(1000, 9999)}" logging.debug(f"[{request_id}] get_chunks called with start_index={start_index}") - + if not self.redis_client: logging.error("Redis not available, cannot retrieve chunks") return [] - + # If no start_index provided, use most recent chunks if start_index is None: start_index = max(0, self.index - 10) # Start closer to current position logging.debug(f"[{request_id}] No start_index provided, using {start_index}") - + # Get current index from Redis current_index = int(self.redis_client.get(self.buffer_index_key) or 0) - + # Calculate range of chunks to retrieve start_id = start_index + 1 chunks_behind = current_index - start_id - + # Adaptive chunk retrieval based on how far behind if chunks_behind > 100: fetch_count = 15 logging.debug(f"[{request_id}] Client very behind ({chunks_behind} chunks), fetching {fetch_count}") elif chunks_behind > 50: - fetch_count = 10 + fetch_count = 10 logging.debug(f"[{request_id}] Client moderately behind ({chunks_behind} chunks), fetching {fetch_count}") elif chunks_behind > 20: fetch_count = 5 @@ -534,94 +555,94 @@ class StreamBuffer: else: fetch_count = 3 logging.debug(f"[{request_id}] Client up-to-date (only {chunks_behind} chunks behind), fetching {fetch_count}") - + end_id = min(current_index + 1, start_id + fetch_count) - + if start_id >= end_id: logging.debug(f"[{request_id}] No new chunks to fetch (start_id={start_id}, end_id={end_id})") return [] - + # Log the range we're retrieving logging.debug(f"[{request_id}] Retrieving chunks {start_id} to {end_id-1} (total: {end_id-start_id})") - + # Directly fetch from Redis using pipeline for efficiency pipe = self.redis_client.pipeline() for idx in range(start_id, end_id): chunk_key = f"{self.buffer_prefix}{idx}" pipe.get(chunk_key) - + results = pipe.execute() - + # Process results chunks = [result for result in results if result is not None] - + # Count non-None results found_chunks = len(chunks) missing_chunks = len(results) - found_chunks - + if missing_chunks > 0: logging.debug(f"[{request_id}] Missing {missing_chunks}/{len(results)} chunks in Redis") - + # Update local tracking if chunks: self.index = end_id - 1 - + # Final log message chunk_sizes = [len(c) for c in chunks] total_bytes = sum(chunk_sizes) if chunks else 0 logging.debug(f"[{request_id}] Returning {len(chunks)} chunks ({total_bytes} bytes)") - + return chunks - + except Exception as e: logging.error(f"Error getting chunks from buffer: {e}", exc_info=True) return [] - + def get_chunks_exact(self, start_index, count): """Get exactly the requested number of chunks from given index""" try: if not self.redis_client: logging.error("Redis not available, cannot retrieve chunks") return [] - + # Calculate range to retrieve start_id = start_index + 1 end_id = start_id + count - + # Get current buffer position current_index = int(self.redis_client.get(self.buffer_index_key) or 0) - + # If requesting beyond current buffer, return what we have if start_id > current_index: return [] - + # Cap end at current buffer position end_id = min(end_id, current_index + 1) - + # Directly fetch from Redis using pipeline pipe = self.redis_client.pipeline() for idx in range(start_id, end_id): chunk_key = f"{self.buffer_prefix}{idx}" pipe.get(chunk_key) - + results = pipe.execute() - + # Filter out None results chunks = [result for result in results if result is not None] - + # Update local index if needed if chunks and start_id + len(chunks) - 1 > self.index: self.index = start_id + len(chunks) - 1 - + return chunks - + except Exception as e: logging.error(f"Error getting exact chunks: {e}", exc_info=True) return [] class ClientManager: """Manages connected clients for a channel with cross-worker visibility""" - + def __init__(self, channel_id, redis_client=None, worker_id=None): self.channel_id = channel_id self.redis_client = redis_client @@ -629,16 +650,16 @@ class ClientManager: self.clients = set() self.lock = threading.Lock() self.last_active_time = time.time() - + # STANDARDIZED KEYS: Move client set under channel namespace self.client_set_key = f"ts_proxy:channel:{channel_id}:clients" self.client_ttl = getattr(Config, 'CLIENT_RECORD_TTL', 60) self.heartbeat_interval = getattr(Config, 'CLIENT_HEARTBEAT_INTERVAL', 10) self.last_heartbeat_time = {} - + # Start heartbeat thread for local clients self._start_heartbeat_thread() - + def _start_heartbeat_thread(self): """Start thread to regularly refresh client presence in Redis""" def heartbeat_task(): @@ -646,20 +667,20 @@ class ClientManager: try: # Wait for the interval time.sleep(self.heartbeat_interval) - + # Send heartbeat for all local clients with self.lock: if not self.clients or not self.redis_client: continue - + # IMPROVED GHOST DETECTION: Check for stale clients before sending heartbeats current_time = time.time() clients_to_remove = set() - + # First identify clients that should be removed for client_id in self.clients: client_key = f"ts_proxy:channel:{self.channel_id}:clients:{client_id}" - + # Check if client exists in Redis at all exists = self.redis_client.exists(client_key) if not exists: @@ -667,123 +688,123 @@ class ClientManager: logging.warning(f"Found ghost client {client_id} - expired in Redis but still in local set") clients_to_remove.add(client_id) continue - + # Check for stale activity using last_active field last_active = self.redis_client.hget(client_key, "last_active") if last_active: last_active_time = float(last_active.decode('utf-8')) time_since_activity = current_time - last_active_time - + # If client hasn't been active for too long, mark for removal # Use configurable threshold for detection ghost_threshold = getattr(Config, 'GHOST_CLIENT_MULTIPLIER', 5.0) if time_since_activity > self.heartbeat_interval * ghost_threshold: logging.warning(f"Detected ghost client {client_id} - last active {time_since_activity:.1f}s ago") clients_to_remove.add(client_id) - + # Remove ghost clients in a separate step for client_id in clients_to_remove: self.remove_client(client_id) - + if clients_to_remove: logging.info(f"Removed {len(clients_to_remove)} ghost clients from channel {self.channel_id}") - + # Now send heartbeats only for remaining clients pipe = self.redis_client.pipeline() current_time = time.time() - + for client_id in self.clients: # Skip clients we just marked for removal if client_id in clients_to_remove: continue - + # Skip if we just sent a heartbeat recently if client_id in self.last_heartbeat_time: time_since_last = current_time - self.last_heartbeat_time[client_id] if time_since_last < self.heartbeat_interval * 0.8: continue - + # Only update clients that remain client_key = f"ts_proxy:channel:{self.channel_id}:clients:{client_id}" pipe.hset(client_key, "last_active", str(current_time)) pipe.expire(client_key, self.client_ttl) - + # Keep client in the set with TTL pipe.sadd(self.client_set_key, client_id) pipe.expire(self.client_set_key, self.client_ttl) - + # Track last heartbeat locally self.last_heartbeat_time[client_id] = current_time - + # Execute all commands atomically pipe.execute() - + # Only notify if we have real clients if self.clients and not all(c in clients_to_remove for c in self.clients): self._notify_owner_of_activity() - + except Exception as e: logging.error(f"Error in client heartbeat thread: {e}") - + thread = threading.Thread(target=heartbeat_task, daemon=True) thread.name = f"client-heartbeat-{self.channel_id}" thread.start() logging.debug(f"Started client heartbeat thread for channel {self.channel_id} (interval: {self.heartbeat_interval}s)") - + def _notify_owner_of_activity(self): """Notify channel owner that clients are active on this worker""" if not self.redis_client or not self.clients: return - + try: worker_id = self.worker_id or "unknown" - + # STANDARDIZED KEY: Worker info under channel namespace worker_key = f"ts_proxy:channel:{self.channel_id}:worker:{worker_id}" self.redis_client.setex(worker_key, self.client_ttl, str(len(self.clients))) - + # STANDARDIZED KEY: Activity timestamp under channel namespace activity_key = f"ts_proxy:channel:{self.channel_id}:activity" self.redis_client.setex(activity_key, self.client_ttl, str(time.time())) except Exception as e: logging.error(f"Error notifying owner of client activity: {e}") - + def add_client(self, client_id, user_agent=None): """Add a client to this channel locally and in Redis""" with self.lock: self.clients.add(client_id) self.last_active_time = time.time() - + if self.redis_client: current_time = str(time.time()) - + # Add to channel's client set self.redis_client.sadd(self.client_set_key, client_id) self.redis_client.expire(self.client_set_key, self.client_ttl) - + # STANDARDIZED KEY: Individual client under channel namespace client_key = f"ts_proxy:channel:{self.channel_id}:clients:{client_id}" - + # Store client info as a hash with all info in one place client_data = { "last_active": current_time, "worker_id": self.worker_id or "unknown", "connect_time": current_time } - + # Add user agent if provided if user_agent: client_data["user_agent"] = user_agent - + # Use HSET to store client data as a hash self.redis_client.hset(client_key, mapping=client_data) self.redis_client.expire(client_key, self.client_ttl) - + # Clear any initialization timer self.redis_client.delete(f"ts_proxy:channel:{self.channel_id}:init_time") - + self._notify_owner_of_activity() - + # Publish client connected event with user agent event_data = { "event": "client_connected", @@ -792,55 +813,55 @@ class ClientManager: "worker_id": self.worker_id or "unknown", "timestamp": time.time() } - + if user_agent: event_data["user_agent"] = user_agent logging.debug(f"Storing user agent '{user_agent}' for client {client_id}") else: logging.debug(f"No user agent provided for client {client_id}") self.redis_client.publish( - f"ts_proxy:events:{self.channel_id}", + f"ts_proxy:events:{self.channel_id}", json.dumps(event_data) ) - + # Get total clients across all workers total_clients = self.get_total_client_count() logging.info(f"New client connected: {client_id} (local: {len(self.clients)}, total: {total_clients})") - + self.last_heartbeat_time[client_id] = time.time() - + return len(self.clients) - + def remove_client(self, client_id): """Remove a client from this channel and Redis""" with self.lock: if client_id in self.clients: self.clients.remove(client_id) - + if client_id in self.last_heartbeat_time: del self.last_heartbeat_time[client_id] - + self.last_active_time = time.time() - + if self.redis_client: # Remove from channel's client set self.redis_client.srem(self.client_set_key, client_id) - + # STANDARDIZED KEY: Delete individual client keys client_key = f"ts_proxy:channel:{self.channel_id}:clients:{client_id}" self.redis_client.delete(client_key) - + # Check if this was the last client remaining = self.redis_client.scard(self.client_set_key) or 0 if remaining == 0: logging.warning(f"Last client removed: {client_id} - channel may shut down soon") - + # Trigger disconnect time tracking even if we're not the owner disconnect_key = f"ts_proxy:channel:{self.channel_id}:last_client_disconnect_time" self.redis_client.setex(disconnect_key, 60, str(time.time())) - + self._notify_owner_of_activity() - + # Publish client disconnected event event_data = json.dumps({ "event": "client_disconnected", @@ -851,41 +872,41 @@ class ClientManager: "remaining_clients": remaining }) self.redis_client.publish(f"ts_proxy:events:{self.channel_id}", event_data) - + total_clients = self.get_total_client_count() logging.info(f"Client disconnected: {client_id} (local: {len(self.clients)}, total: {total_clients})") - + return len(self.clients) - + def get_client_count(self): """Get local client count""" with self.lock: return len(self.clients) - + def get_total_client_count(self): """Get total client count across all workers""" if not self.redis_client: return len(self.clients) - + try: # Count members in the client set return self.redis_client.scard(self.client_set_key) or 0 except Exception as e: logging.error(f"Error getting total client count: {e}") return len(self.clients) # Fall back to local count - + def refresh_client_ttl(self): """Refresh TTL for active clients to prevent expiration""" if not self.redis_client: return - + try: # Refresh TTL for all clients belonging to this worker for client_id in self.clients: # STANDARDIZED: Use channel namespace for client keys client_key = f"ts_proxy:channel:{self.channel_id}:clients:{client_id}" self.redis_client.expire(client_key, self.client_ttl) - + # Refresh TTL on the set itself self.redis_client.expire(self.client_set_key, self.client_ttl) except Exception as e: @@ -893,7 +914,7 @@ class ClientManager: class StreamFetcher: """Handles stream data fetching""" - + def __init__(self, manager: StreamManager, buffer: StreamBuffer): self.manager = manager self.buffer = buffer @@ -919,10 +940,10 @@ class StreamFetcher: if not self.manager.should_retry(): logging.error(f"Failed to connect after {self.manager.max_retries} attempts") return False - + if not self.manager.running: return False - + self.manager.retry_count += 1 logging.info(f"Connecting to stream: {self.manager.url} " f"(attempt {self.manager.retry_count}/{self.manager.max_retries})") @@ -941,13 +962,13 @@ class StreamFetcher: if not self.manager.running: logging.info("Stream fetch stopped - shutting down") return - + if chunk: if self.manager.ready_event.is_set(): logging.info("Stream switch in progress, closing connection") self.manager.ready_event.clear() break - + with self.buffer.lock: self.buffer.buffer.append(chunk) self.buffer.index += 1 @@ -956,10 +977,10 @@ class StreamFetcher: """Handle stream connection errors""" logging.error(f"Stream connection error: {error}") self.manager.connected = False - + if not self.manager.running: return - + logging.info(f"Attempting to reconnect in {Config.RECONNECT_DELAY} seconds...") if not wait_for_running(self.manager, Config.RECONNECT_DELAY): return @@ -975,26 +996,26 @@ def wait_for_running(manager: StreamManager, delay: float) -> bool: class ProxyServer: """Manages TS proxy server instance with worker coordination""" - + def __init__(self): """Initialize proxy server with worker identification""" self.stream_managers = {} self.stream_buffers = {} self.client_managers = {} - + # Generate a unique worker ID import socket import os pid = os.getpid() hostname = socket.gethostname() self.worker_id = f"{hostname}:{pid}" - + # Connect to Redis self.redis_client = None try: import redis from django.conf import settings - + redis_url = getattr(settings, 'REDIS_URL', 'redis://localhost:6379/0') self.redis_client = redis.from_url(redis_url) logging.info(f"Connected to Redis at {redis_url}") @@ -1002,37 +1023,37 @@ class ProxyServer: except Exception as e: self.redis_client = None logging.error(f"Failed to connect to Redis: {e}") - + # Start cleanup thread self.cleanup_interval = getattr(Config, 'CLEANUP_INTERVAL', 60) self._start_cleanup_thread() - + # Start event listener for Redis pubsub messages self._start_event_listener() - + def _start_event_listener(self): """Listen for events from other workers""" if not self.redis_client: return - + def event_listener(): try: pubsub = self.redis_client.pubsub() pubsub.psubscribe("ts_proxy:events:*") - + logging.info("Started Redis event listener for client activity") - + for message in pubsub.listen(): if message["type"] != "pmessage": continue - + try: channel = message["channel"].decode("utf-8") data = json.loads(message["data"].decode("utf-8")) - + event_type = data.get("event") channel_id = data.get("channel_id") - + if channel_id and event_type: # For owner, update client status immediately if self.am_i_owner(channel_id): @@ -1042,7 +1063,7 @@ class ProxyServer: # RENAMED: no_clients_since → last_client_disconnect_time disconnect_key = f"ts_proxy:channel:{channel_id}:last_client_disconnect_time" self.redis_client.delete(disconnect_key) - + elif event_type == "client_disconnected": logging.debug(f"Owner received client_disconnected event for channel {channel_id}") # Check if any clients remain @@ -1050,27 +1071,27 @@ class ProxyServer: # VERIFY REDIS CLIENT COUNT DIRECTLY client_set_key = f"ts_proxy:channel:{channel_id}:clients" total = self.redis_client.scard(client_set_key) or 0 - + if total == 0: logging.debug(f"No clients left after disconnect event - stopping channel {channel_id}") # Set the disconnect timer for other workers to see disconnect_key = f"ts_proxy:channel:{channel_id}:last_client_disconnect_time" self.redis_client.setex(disconnect_key, 60, str(time.time())) - + # Get configured shutdown delay or default shutdown_delay = getattr(Config, 'CHANNEL_SHUTDOWN_DELAY', 0) - + if shutdown_delay > 0: logging.info(f"Waiting {shutdown_delay}s before stopping channel...") time.sleep(shutdown_delay) - + # Re-check client count before stopping total = self.redis_client.scard(client_set_key) or 0 if total > 0: logging.info(f"New clients connected during shutdown delay - aborting shutdown") self.redis_client.delete(disconnect_key) return - + # Stop the channel directly self.stop_channel(channel_id) @@ -1080,7 +1101,7 @@ class ProxyServer: # Handle stream switch request new_url = data.get("url") user_agent = data.get("user_agent") - + if new_url and channel_id in self.stream_managers: # Update metadata in Redis if self.redis_client: @@ -1088,18 +1109,18 @@ class ProxyServer: self.redis_client.hset(metadata_key, "url", new_url) if user_agent: self.redis_client.hset(metadata_key, "user_agent", user_agent) - + # Set switch status status_key = f"ts_proxy:channel:{channel_id}:switch_status" self.redis_client.set(status_key, "switching") - + # Perform the stream switch stream_manager = self.stream_managers[channel_id] success = stream_manager.update_url(new_url) - + if success: logging.info(f"Stream switch initiated for channel {channel_id}") - + # Publish confirmation switch_result = { "event": "stream_switched", @@ -1109,16 +1130,16 @@ class ProxyServer: "timestamp": time.time() } self.redis_client.publish( - f"ts_proxy:events:{channel_id}", + f"ts_proxy:events:{channel_id}", json.dumps(switch_result) ) - + # Update status if self.redis_client: self.redis_client.set(status_key, "switched") else: logging.error(f"Failed to switch stream for channel {channel_id}") - + # Publish failure switch_result = { "event": "stream_switched", @@ -1128,7 +1149,7 @@ class ProxyServer: "timestamp": time.time() } self.redis_client.publish( - f"ts_proxy:events:{channel_id}", + f"ts_proxy:events:{channel_id}", json.dumps(switch_result) ) except Exception as e: @@ -1138,7 +1159,7 @@ class ProxyServer: time.sleep(5) # Wait before reconnecting # Try to restart the listener self._start_event_listener() - + thread = threading.Thread(target=event_listener, daemon=True) thread.name = "redis-event-listener" thread.start() @@ -1147,7 +1168,7 @@ class ProxyServer: """Get the worker ID that owns this channel with proper error handling""" if not self.redis_client: return None - + try: lock_key = f"ts_proxy:channel:{channel_id}:owner" owner = self.redis_client.get(lock_key) @@ -1157,30 +1178,30 @@ class ProxyServer: except Exception as e: logging.error(f"Error getting channel owner: {e}") return None - + def am_i_owner(self, channel_id): """Check if this worker is the owner of the channel""" owner = self.get_channel_owner(channel_id) return owner == self.worker_id - + def try_acquire_ownership(self, channel_id, ttl=30): """Try to become the owner of this channel using proper locking""" if not self.redis_client: return True # If no Redis, always become owner - + try: # Create a lock key with proper namespace lock_key = f"ts_proxy:channel:{channel_id}:owner" - + # Use Redis SETNX for atomic locking - only succeeds if the key doesn't exist acquired = self.redis_client.setnx(lock_key, self.worker_id) - + # If acquired, set expiry to prevent orphaned locks if acquired: self.redis_client.expire(lock_key, ttl) logging.info(f"Worker {self.worker_id} acquired ownership of channel {channel_id}") return True - + # If not acquired, check if we already own it (might be a retry) current_owner = self.redis_client.get(lock_key) if current_owner and current_owner.decode('utf-8') == self.worker_id: @@ -1188,22 +1209,22 @@ class ProxyServer: self.redis_client.expire(lock_key, ttl) logging.info(f"Worker {self.worker_id} refreshed ownership of channel {channel_id}") return True - + # Someone else owns it return False - + except Exception as e: logging.error(f"Error acquiring channel ownership: {e}") return False - + def release_ownership(self, channel_id): """Release ownership of this channel safely""" if not self.redis_client: return - + try: lock_key = f"ts_proxy:channel:{channel_id}:owner" - + # Only delete if we're the current owner to prevent race conditions current = self.redis_client.get(lock_key) if current and current.decode('utf-8') == self.worker_id: @@ -1211,16 +1232,16 @@ class ProxyServer: logging.info(f"Released ownership of channel {channel_id}") except Exception as e: logging.error(f"Error releasing channel ownership: {e}") - + def extend_ownership(self, channel_id, ttl=30): """Extend ownership lease with grace period""" if not self.redis_client: return False - + try: - lock_key = f"ts_proxy:channel:{channel_id}:owner" + lock_key = f"ts_proxy:channel:{channel_id}:owner" current = self.redis_client.get(lock_key) - + # Only extend if we're still the owner if current and current.decode('utf-8') == self.worker_id: self.redis_client.expire(lock_key, ttl) @@ -1229,73 +1250,73 @@ class ProxyServer: except Exception as e: logging.error(f"Error extending ownership: {e}") return False - - def initialize_channel(self, url, channel_id, user_agent=None): + + def initialize_channel(self, url, channel_id, user_agent=None, transcode_cmd=None): """Initialize a channel without redundant active key""" try: # Get channel URL from Redis if available channel_url = url channel_user_agent = user_agent - + # First check if channel metadata already exists existing_metadata = None metadata_key = f"ts_proxy:channel:{channel_id}:metadata" - + if self.redis_client: existing_metadata = self.redis_client.hgetall(metadata_key) - + # If no url was passed, try to get from Redis if not url and existing_metadata: url_bytes = existing_metadata.get(b'url') if url_bytes: channel_url = url_bytes.decode('utf-8') - + ua_bytes = existing_metadata.get(b'user_agent') if ua_bytes: channel_user_agent = ua_bytes.decode('utf-8') - + # Check if channel is already owned current_owner = self.get_channel_owner(channel_id) - + # Exit early if another worker owns the channel if current_owner and current_owner != self.worker_id: logging.info(f"Channel {channel_id} already owned by worker {current_owner}") logging.info(f"This worker ({self.worker_id}) will read from Redis buffer only") - + # Create buffer but not stream manager buffer = StreamBuffer(channel_id=channel_id, redis_client=self.redis_client) self.stream_buffers[channel_id] = buffer - + # Create client manager with channel_id and redis_client client_manager = ClientManager(channel_id=channel_id, redis_client=self.redis_client, worker_id=self.worker_id) self.client_managers[channel_id] = client_manager - + return True - + # Only continue with full initialization if URL is provided # or we can get it from Redis if not channel_url: logging.error(f"No URL available for channel {channel_id}") return False - + # Try to acquire ownership with Redis locking if not self.try_acquire_ownership(channel_id): # Another worker just acquired ownership logging.info(f"Another worker just acquired ownership of channel {channel_id}") - + # Create buffer but not stream manager buffer = StreamBuffer(channel_id=channel_id, redis_client=self.redis_client) self.stream_buffers[channel_id] = buffer - + # Create client manager with channel_id and redis_client client_manager = ClientManager(channel_id=channel_id, redis_client=self.redis_client, worker_id=self.worker_id) self.client_managers[channel_id] = client_manager - + return True - + # We now own the channel - ONLY NOW should we set metadata with initializing state logging.info(f"Worker {self.worker_id} is now the owner of channel {channel_id}") - + if self.redis_client: # NOW create or update metadata with initializing state metadata = { @@ -1307,49 +1328,49 @@ class ProxyServer: } if channel_user_agent: metadata["user_agent"] = channel_user_agent - + # Set channel metadata self.redis_client.hset(metadata_key, mapping=metadata) self.redis_client.expire(metadata_key, 3600) # 1 hour TTL - + # Create stream buffer buffer = StreamBuffer(channel_id=channel_id, redis_client=self.redis_client) logging.debug(f"Created StreamBuffer for channel {channel_id}") 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) + stream_manager = StreamManager(channel_url, buffer, user_agent=channel_user_agent, transcode_cmd=transcode_cmd) logging.debug(f"Created StreamManager for channel {channel_id}") self.stream_managers[channel_id] = stream_manager - + # Create client manager with channel_id, redis_client AND worker_id client_manager = ClientManager( - channel_id=channel_id, + channel_id=channel_id, redis_client=self.redis_client, worker_id=self.worker_id ) self.client_managers[channel_id] = client_manager - + # Start stream manager thread only for the owner thread = threading.Thread(target=stream_manager.run, daemon=True) thread.name = f"stream-{channel_id}" thread.start() logging.info(f"Started stream manager thread for channel {channel_id}") - + # If we're the owner, we need to set the channel state rather than starting a grace period immediately if self.am_i_owner(channel_id): self.update_channel_state(channel_id, "connecting", { "init_time": str(time.time()), "owner": self.worker_id }) - + # Set connection attempt start time attempt_key = f"ts_proxy:channel:{channel_id}:connection_attempt_time" self.redis_client.setex(attempt_key, 60, str(time.time())) - - logging.info(f"Channel {channel_id} in connecting state - will start grace period after connection") + + logging.info(f"Channel {channel_id} in connecting state - will start grace period after connection") return True - + except Exception as e: logging.error(f"Error initializing channel {channel_id}: {e}", exc_info=True) # Release ownership on failure @@ -1361,40 +1382,40 @@ class ProxyServer: # Check local memory first if channel_id in self.stream_managers or channel_id in self.stream_buffers: return True - + # Check Redis using the standard key pattern if self.redis_client: # Primary check - look for channel metadata metadata_key = f"ts_proxy:channel:{channel_id}:metadata" - + # If metadata exists, return true if self.redis_client.exists(metadata_key): return True - + # Additional checks if metadata doesn't exist additional_keys = [ f"ts_proxy:channel:{channel_id}:clients", f"ts_proxy:channel:{channel_id}:buffer:index", f"ts_proxy:channel:{channel_id}:owner" ] - + for key in additional_keys: if self.redis_client.exists(key): return True - + return False def stop_channel(self, channel_id): """Stop a channel with proper ownership handling""" try: logging.info(f"Stopping channel {channel_id}") - + # Only stop the actual stream manager if we're the owner if self.am_i_owner(channel_id): logging.info(f"This worker ({self.worker_id}) is the owner - closing provider connection") if channel_id in self.stream_managers: stream_manager = self.stream_managers[channel_id] - + # Signal thread to stop and close resources if hasattr(stream_manager, 'stop'): stream_manager.stop() @@ -1402,16 +1423,16 @@ class ProxyServer: stream_manager.running = False if hasattr(stream_manager, '_close_socket'): stream_manager._close_socket() - + # Wait for stream thread to finish stream_thread_name = f"stream-{channel_id}" stream_thread = None - + for thread in threading.enumerate(): if thread.name == stream_thread_name: stream_thread = thread break - + if stream_thread and stream_thread.is_alive(): logging.info(f"Waiting for stream thread to terminate") try: @@ -1421,27 +1442,27 @@ class ProxyServer: logging.warning(f"Stream thread did not terminate within timeout") except RuntimeError: logging.debug("Could not join stream thread (may be current thread)") - + # Release ownership self.release_ownership(channel_id) logging.info(f"Released ownership of channel {channel_id}") - + # Always clean up local resources if channel_id in self.stream_managers: del self.stream_managers[channel_id] logging.info(f"Removed stream manager for channel {channel_id}") - + if channel_id in self.stream_buffers: del self.stream_buffers[channel_id] logging.info(f"Removed stream buffer for channel {channel_id}") - + if channel_id in self.client_managers: del self.client_managers[channel_id] logging.info(f"Removed client manager for channel {channel_id}") - + # Clean up Redis keys self._clean_redis_keys(channel_id) - + return True except Exception as e: logging.error(f"Error stopping channel {channel_id}: {e}", exc_info=True) @@ -1450,18 +1471,18 @@ class ProxyServer: def check_inactive_channels(self): """Check for inactive channels (no clients) and stop them""" channels_to_stop = [] - + for channel_id, client_manager in self.client_managers.items(): if client_manager.get_client_count() == 0: channels_to_stop.append(channel_id) - + for channel_id in channels_to_stop: logging.info(f"Auto-stopping inactive channel {channel_id}") self.stop_channel(channel_id) def _cleanup_channel(self, channel_id: str) -> None: """Remove channel resources""" - for collection in [self.stream_managers, self.stream_buffers, + for collection in [self.stream_managers, self.stream_buffers, self.client_managers, self.fetch_threads]: collection.pop(channel_id, None) @@ -1480,7 +1501,7 @@ class ProxyServer: if self.am_i_owner(channel_id): # Extend ownership lease self.extend_ownership(channel_id) - + # Get channel state from metadata hash channel_state = "unknown" if self.redis_client: @@ -1494,11 +1515,11 @@ class ProxyServer: if channel_id in self.client_managers: client_manager = self.client_managers[channel_id] total_clients = client_manager.get_total_client_count() - + # Log client count periodically if time.time() % 30 < 1: # Every ~30 seconds logging.info(f"Channel {channel_id} has {total_clients} clients, state: {channel_state}") - + # If in connecting or waiting_for_clients state, check grace period if channel_state in ["connecting", "waiting_for_clients"]: # Get connection ready time from metadata @@ -1508,22 +1529,22 @@ class ProxyServer: connection_ready_time = float(metadata[b'connection_ready_time'].decode('utf-8')) except (ValueError, TypeError): pass - + # If still connecting, give it more time if channel_state == "connecting": logging.debug(f"Channel {channel_id} still connecting - not checking for clients yet") continue - + # If waiting for clients, check grace period if connection_ready_time: grace_period = getattr(Config, 'CHANNEL_INIT_GRACE_PERIOD', 20) time_since_ready = time.time() - connection_ready_time - + # Add this debug log - logging.debug(f"GRACE PERIOD CHECK: Channel {channel_id} in {channel_state} state, " + logging.debug(f"GRACE PERIOD CHECK: Channel {channel_id} in {channel_state} state, " f"time_since_ready={time_since_ready:.1f}s, grace_period={grace_period}s, " f"total_clients={total_clients}") - + if time_since_ready <= grace_period: # Still within grace period logging.debug(f"Channel {channel_id} in grace period - {time_since_ready:.1f}s of {grace_period}s elapsed") @@ -1548,7 +1569,7 @@ class ProxyServer: # Check if there's a pending no-clients timeout disconnect_key = f"ts_proxy:channel:{channel_id}:last_client_disconnect_time" disconnect_time = None - + if self.redis_client: disconnect_value = self.redis_client.get(disconnect_key) if disconnect_value: @@ -1556,9 +1577,9 @@ class ProxyServer: disconnect_time = float(disconnect_value.decode('utf-8')) except (ValueError, TypeError) as e: logging.error(f"Invalid disconnect time for channel {channel_id}: {e}") - + current_time = time.time() - + if not disconnect_time: # First time seeing zero clients, set timestamp if self.redis_client: @@ -1570,19 +1591,19 @@ class ProxyServer: self.stop_channel(channel_id) else: # Still in shutdown delay period - logging.debug(f"Channel {channel_id} shutdown timer: " - f"{current_time - disconnect_time:.1f}s of " + logging.debug(f"Channel {channel_id} shutdown timer: " + f"{current_time - disconnect_time:.1f}s of " f"{getattr(Config, 'CHANNEL_SHUTDOWN_DELAY', 5)}s elapsed") else: # There are clients or we're still connecting - clear any disconnect timestamp if self.redis_client: self.redis_client.delete(f"ts_proxy:channel:{channel_id}:last_client_disconnect_time") - + except Exception as e: logging.error(f"Error in cleanup thread: {e}", exc_info=True) - + time.sleep(getattr(Config, 'CLEANUP_CHECK_INTERVAL', 1)) - + thread = threading.Thread(target=cleanup_task, daemon=True) thread.name = "ts-proxy-cleanup" thread.start() @@ -1592,28 +1613,28 @@ class ProxyServer: """Check for orphaned channels in Redis (owner worker crashed)""" if not self.redis_client: return - + try: # Get all active channel keys channel_pattern = "ts_proxy:channel:*:metadata" channel_keys = self.redis_client.keys(channel_pattern) - + for key in channel_keys: try: channel_id = key.decode('utf-8').split(':')[2] - + # Skip channels we already have locally if channel_id in self.stream_buffers: continue - + # Check if this channel has an owner owner = self.get_channel_owner(channel_id) - + if not owner: # Check if there are any clients client_set_key = f"ts_proxy:channel:{channel_id}:clients" client_count = self.redis_client.scard(client_set_key) or 0 - + if client_count > 0: # Orphaned channel with clients - we could take ownership logging.info(f"Found orphaned channel {channel_id} with {client_count} clients") @@ -1623,7 +1644,7 @@ class ProxyServer: self._clean_redis_keys(channel_id) except Exception as e: logging.error(f"Error processing channel key {key}: {e}") - + except Exception as e: logging.error(f"Error checking orphaned channels: {e}") @@ -1631,16 +1652,18 @@ class ProxyServer: """Clean up all Redis keys for a channel""" if not self.redis_client: return - + try: + channel = Channel.objects.get(id=channel_id) + channel.release_stream() # All keys are now under the channel namespace for easy pattern matching channel_pattern = f"ts_proxy:channel:{channel_id}:*" all_keys = self.redis_client.keys(channel_pattern) - + if all_keys: self.redis_client.delete(*all_keys) logging.info(f"Cleaned up {len(all_keys)} Redis keys for channel {channel_id}") - + except Exception as e: logging.error(f"Error cleaning Redis keys for channel {channel_id}: {e}") @@ -1648,12 +1671,12 @@ class ProxyServer: """Refresh TTL for active channels using standard keys""" if not self.redis_client: return - + # Refresh registry entries for channels we own for channel_id in self.stream_managers.keys(): # Use standard key pattern metadata_key = f"ts_proxy:channel:{channel_id}:metadata" - + # Update activity timestamp in metadata only self.redis_client.hset(metadata_key, "last_active", str(time.time())) self.redis_client.expire(metadata_key, 3600) # Reset TTL on metadata hash @@ -1662,38 +1685,37 @@ class ProxyServer: """Update channel state with proper history tracking and logging""" if not self.redis_client: return False - + try: metadata_key = f"ts_proxy:channel:{channel_id}:metadata" - + # Get current state for logging current_state = None metadata = self.redis_client.hgetall(metadata_key) if metadata and b'state' in metadata: current_state = metadata[b'state'].decode('utf-8') - + # Only update if state is actually changing if current_state == new_state: logging.debug(f"Channel {channel_id} state unchanged: {current_state}") return True - + # Prepare update data update_data = { "state": new_state, "state_changed_at": str(time.time()) } - + # Add optional additional fields if additional_fields: update_data.update(additional_fields) - + # Update the metadata self.redis_client.hset(metadata_key, mapping=update_data) - + # Log the transition logging.info(f"Channel {channel_id} state transition: {current_state or 'None'} → {new_state}") return True except Exception as e: logging.error(f"Error updating channel state: {e}") return False - diff --git a/apps/proxy/ts_proxy/urls.py b/apps/proxy/ts_proxy/urls.py index 77c22534..d53ade8d 100644 --- a/apps/proxy/ts_proxy/urls.py +++ b/apps/proxy/ts_proxy/urls.py @@ -5,6 +5,6 @@ app_name = 'ts_proxy' urlpatterns = [ path('stream/', views.stream_ts, name='stream'), - path('initialize/', views.initialize_stream, name='initialize'), + # path('initialize/', views.initialize_stream, name='initialize'), path('change_stream/', views.change_stream, name='change_stream'), -] \ No newline at end of file +] diff --git a/apps/proxy/ts_proxy/views.py b/apps/proxy/ts_proxy/views.py index 23b41e1c..5516d20e 100644 --- a/apps/proxy/ts_proxy/views.py +++ b/apps/proxy/ts_proxy/views.py @@ -5,11 +5,16 @@ import time import random import sys import os +import re from django.http import StreamingHttpResponse, JsonResponse from django.views.decorators.csrf import csrf_exempt from django.views.decorators.http import require_http_methods, require_GET +from django.shortcuts import get_object_or_404 from apps.proxy.config import TSConfig as Config from .server import ProxyServer +from apps.channels.models import Channel, Stream +from apps.m3u.models import M3UAccount, M3UAccountProfile +from core.models import UserAgent # Configure logging properly to ensure visibility logger = logging.getLogger(__name__) @@ -26,24 +31,14 @@ print("TS PROXY VIEWS INITIALIZED", file=sys.stderr) # Initialize proxy server proxy_server = ProxyServer() -@csrf_exempt -@require_http_methods(["POST"]) -def initialize_stream(request, channel_id): +def initialize_stream(channel_id, url, user_agent, transcode_cmd): """Initialize a new stream channel with initialization-based ownership""" try: - data = json.loads(request.body) - url = data.get('url') - if not url: - return JsonResponse({'error': 'No URL provided'}, status=400) - - # Get optional user_agent from request - user_agent = data.get('user_agent') - # Try to acquire ownership and create connection - success = proxy_server.initialize_channel(url, channel_id, user_agent) + success = proxy_server.initialize_channel(url, channel_id, user_agent, transcode_cmd) if not success: return JsonResponse({'error': 'Failed to initialize channel'}, status=500) - + # If we're the owner, wait for connection if proxy_server.am_i_owner(channel_id): # Wait for connection to be established @@ -62,7 +57,7 @@ def initialize_stream(request, channel_id): 'error': 'Failed to connect' }, status=502) time.sleep(0.1) - + # Return success response with owner status return JsonResponse({ 'message': 'Stream initialized and connected', @@ -70,7 +65,7 @@ def initialize_stream(request, channel_id): 'url': url, 'owner': proxy_server.am_i_owner(channel_id) }) - + except json.JSONDecodeError: return JsonResponse({'error': 'Invalid JSON'}, status=400) except Exception as e: @@ -80,24 +75,54 @@ def initialize_stream(request, channel_id): @require_GET def stream_ts(request, channel_id): """Stream TS data to client with improved waiting for initialization""" + user_agent = None + channel = get_object_or_404(Channel, pk=channel_id) + try: # Check if channel exists or initialize it if not proxy_server.check_if_channel_exists(channel_id): - return JsonResponse({'error': 'Channel not found'}, status=404) + stream_id, profile_id = channel.get_stream() + if stream_id is None or profile_id is None: + return JsonResponse({'error': 'Channel not available'}, status=404) + + # Load in necessary objects for the stream + stream = get_object_or_404(Stream, pk=stream_id) + profile = get_object_or_404(M3UAccountProfile, pk=profile_id) + + # Load in the user-agent for the account + m3u_account = M3UAccount.objects.get(id=profile.m3u_account.id) + user_agent = UserAgent.objects.get(id=m3u_account.user_agent.id).user_agent + + # Generate stream URL based on the selected profile + input_url = stream.custom_url or stream.url + logger.debug("Executing the following pattern replacement:") + logger.debug(f" search: {profile.search_pattern}") + safe_replace_pattern = re.sub(r'\$(\d+)', r'\\\1', profile.replace_pattern) + logger.debug(f" replace: {profile.replace_pattern}") + logger.debug(f" safe replace: {safe_replace_pattern}") + stream_url = re.sub(profile.search_pattern, safe_replace_pattern, input_url) + logger.debug(f"Generated stream url: {stream_url}") + + # Generate transcode command + # @TODO: once complete, provide option to direct proxy + transcode_cmd = channel.get_stream_profile().build_command(stream_url) + + initialize_stream(channel_id, stream_url, user_agent, transcode_cmd) + # Get user agent from request headers - user_agent = None - for header in ['HTTP_USER_AGENT', 'User-Agent', 'user-agent']: - if header in request.META: - user_agent = request.META[header] - logger.debug(f"Found user agent in header: {header}") - break - + if user_agent is None: + for header in ['HTTP_USER_AGENT', 'User-Agent', 'user-agent']: + if header in request.META: + user_agent = request.META[header] + logger.debug(f"Found user agent in header: {header}") + break + # Wait for channel to become ready if it's initializing if proxy_server.redis_client: wait_start = time.time() max_wait = getattr(Config, 'CLIENT_WAIT_TIMEOUT', 30) # Maximum wait time in seconds - + # Check channel state metadata_key = f"ts_proxy:channel:{channel_id}:metadata" while time.time() - wait_start < max_wait: @@ -105,9 +130,9 @@ def stream_ts(request, channel_id): if not metadata or b'state' not in metadata: logger.warning(f"Channel {channel_id} metadata missing") break - + state = metadata[b'state'].decode('utf-8') - + # If channel is already active or waiting for clients, no need to wait if state in ['waiting_for_clients', 'active']: logger.debug(f"Channel {channel_id} ready (state={state})") @@ -121,7 +146,7 @@ def stream_ts(request, channel_id): # Unknown or error state logger.warning(f"Channel {channel_id} in unexpected state: {state}") break - + # Check if we timed out waiting if time.time() - wait_start >= max_wait: logger.warning(f"Timeout waiting for channel {channel_id} to become ready") @@ -131,7 +156,7 @@ def stream_ts(request, channel_id): # This handles the case where channel exists in Redis but not in this worker if channel_id not in proxy_server.stream_buffers or channel_id not in proxy_server.client_managers: logger.warning(f"Channel {channel_id} exists in Redis but not initialized in this worker - initializing now") - + # Get URL from Redis metadata if available url = None metadata_key = f"ts_proxy:channel:{channel_id}:metadata" @@ -139,15 +164,15 @@ def stream_ts(request, channel_id): url_bytes = proxy_server.redis_client.hget(metadata_key, "url") if url_bytes: url = url_bytes.decode('utf-8') - + # Initialize local resources (won't recreate stream connection if another worker owns it) success = proxy_server.initialize_channel(url, channel_id, user_agent) if not success: logger.error(f"Failed to initialize channel {channel_id} locally") return JsonResponse({'error': 'Failed to initialize channel locally'}, status=500) - + logger.info(f"Successfully initialized channel {channel_id} locally") - + # Double-check after initialization if channel_id not in proxy_server.client_managers: logger.error(f"Critical error: Channel {channel_id} client manager still missing after initialization") @@ -159,20 +184,20 @@ def stream_ts(request, channel_id): stream_start_time = time.time() bytes_sent = 0 chunks_sent = 0 - + try: # ENHANCED USER AGENT DETECTION - check multiple possible headers user_agent = None - + # Try multiple possible header formats ua_headers = ['HTTP_USER_AGENT', 'User-Agent', 'user-agent', 'User_Agent'] - + for header in ua_headers: if header in request.META: user_agent = request.META[header] logger.debug(f"Found user agent in header: {header}") break - + # Try request.headers dictionary (Django 2.2+) if not user_agent and hasattr(request, 'headers'): for header in ['User-Agent', 'user-agent']: @@ -180,7 +205,7 @@ def stream_ts(request, channel_id): user_agent = request.headers[header] logger.debug(f"Found user agent in request.headers: {header}") break - + # Final fallback - check if in any header with case-insensitive matching if not user_agent: for key, value in request.META.items(): @@ -188,63 +213,63 @@ def stream_ts(request, channel_id): user_agent = value logger.debug(f"Found user agent in alternate header: {key}") break - + # Log headers for debugging user agent issues if not user_agent: # Log all headers to help troubleshoot headers = {k: v for k, v in request.META.items() if k.startswith('HTTP_')} logger.debug(f"No user agent found in request. Available headers: {headers}") user_agent = "Unknown-Client" # Default value instead of None - + logger.info(f"[{client_id}] New client connected to channel {channel_id} with user agent: {user_agent}") - + # Add client to manager with user agent client_manager = proxy_server.client_managers[channel_id] client_count = client_manager.add_client(client_id, user_agent) - + # If this is the first client, try to acquire ownership if client_count == 1 and not proxy_server.am_i_owner(channel_id): if proxy_server.try_acquire_ownership(channel_id): logger.info(f"[{client_id}] First client, acquiring channel ownership") - + # Get channel metadata from Redis if proxy_server.redis_client: 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") - + url = url_bytes.decode('utf-8') if url_bytes else None user_agent = ua_bytes.decode('utf-8') if ua_bytes else None - + if url: # Create and start stream connection from .server import StreamManager # Import here to avoid circular import - + logger.info(f"[{client_id}] Creating stream connection for URL: {url}") buffer = proxy_server.stream_buffers[channel_id] - + stream_manager = StreamManager(url, buffer, user_agent=user_agent) proxy_server.stream_managers[channel_id] = stream_manager - + thread = threading.Thread(target=stream_manager.run, daemon=True) thread.name = f"stream-{channel_id}" thread.start() - + # Wait briefly for connection wait_start = time.time() while not stream_manager.connected: if time.time() - wait_start > Config.CONNECTION_TIMEOUT: break time.sleep(0.1) - + # Get buffer - stream manager may not exist in this worker buffer = proxy_server.stream_buffers.get(channel_id) stream_manager = proxy_server.stream_managers.get(channel_id) - + if not buffer: logger.error(f"[{client_id}] No buffer found for channel {channel_id}") return - + # Client state tracking - use config for initial position local_index = max(0, buffer.index - Config.INITIAL_BEHIND_CHUNKS) initial_position = local_index @@ -254,46 +279,46 @@ def stream_ts(request, channel_id): chunks_sent = 0 stream_start_time = time.time() consecutive_empty = 0 # Track consecutive empty reads - + # Timing parameters from config ts_packet_size = 188 - target_bitrate = Config.TARGET_BITRATE + target_bitrate = Config.TARGET_BITRATE packets_per_second = target_bitrate / (8 * ts_packet_size) - + logger.info(f"[{client_id}] Starting stream at index {local_index} (buffer at {buffer.index})") - + # Check if we're the owner worker is_owner_worker = proxy_server.am_i_owner(channel_id) if hasattr(proxy_server, 'am_i_owner') else True - + # Main streaming loop while True: # Get chunks at client's position chunks = buffer.get_chunks_exact(local_index, Config.CHUNK_BATCH_SIZE) - + if chunks: # Reset empty counters since we got data empty_reads = 0 consecutive_empty = 0 - + # Track and send chunks chunk_sizes = [len(c) for c in chunks] total_size = sum(chunk_sizes) start_idx = local_index + 1 end_idx = local_index + len(chunks) - + logger.debug(f"[{client_id}] Retrieved {len(chunks)} chunks ({total_size} bytes) from index {start_idx} to {end_idx}") - + # Calculate total packet count for this batch to maintain timing total_packets = sum(len(chunk) // ts_packet_size for chunk in chunks) batch_start_time = time.time() packets_sent_in_batch = 0 - + # Send chunks with pacing for chunk in chunks: packets_in_chunk = len(chunk) // ts_packet_size bytes_sent += len(chunk) chunks_sent += 1 - + # CRITICAL FIX: Detect client disconnection when yielding data try: yield chunk @@ -303,24 +328,24 @@ def stream_ts(request, channel_id): if channel_id in proxy_server.client_managers: proxy_server.client_managers[channel_id].remove_client(client_id) break # Exit the generator - + # Pacing logic packets_sent_in_batch += packets_in_chunk elapsed = time.time() - batch_start_time target_time = packets_sent_in_batch / packets_per_second - + # If we're sending too fast, add a small delay if elapsed < target_time and packets_sent_in_batch < total_packets: sleep_time = min(target_time - elapsed, 0.05) if sleep_time > 0.001: time.sleep(sleep_time) - + # Log progress periodically if chunks_sent % 100 == 0: elapsed = time.time() - stream_start_time rate = bytes_sent / elapsed / 1024 if elapsed > 0 else 0 logger.info(f"[{client_id}] Stats: {chunks_sent} chunks, {bytes_sent/1024:.1f}KB, {rate:.1f}KB/s") - + # Update local index local_index = end_idx last_yield_time = time.time() @@ -328,14 +353,14 @@ def stream_ts(request, channel_id): # No chunks available empty_reads += 1 consecutive_empty += 1 - + # Check if we're caught up to buffer head at_buffer_head = local_index >= buffer.index - + # If we're at buffer head and no data is coming, send keepalive # Only check stream manager health if it exists stream_healthy = stream_manager.healthy if stream_manager else True - + if at_buffer_head and not stream_healthy and consecutive_empty >= 5: # Create a null TS packet as keepalive (188 bytes filled with padding) # This prevents VLC from hitting EOF @@ -343,7 +368,7 @@ def stream_ts(request, channel_id): keepalive_packet[0] = 0x47 # Sync byte keepalive_packet[1] = 0x1F # PID high bits (null packet) keepalive_packet[2] = 0xFF # PID low bits (null packet) - + logger.debug(f"[{client_id}] Sending keepalive packet while waiting at buffer head") yield bytes(keepalive_packet) bytes_sent += len(keepalive_packet) @@ -354,12 +379,12 @@ def stream_ts(request, channel_id): # Standard wait sleep_time = min(0.1 * consecutive_empty, 1.0) # Progressive backoff up to 1s time.sleep(sleep_time) - + # Log empty reads periodically if empty_reads % 50 == 0: stream_status = "healthy" if (stream_manager and stream_manager.healthy) else "unknown" logger.debug(f"[{client_id}] Waiting for chunks beyond {local_index} (buffer at {buffer.index}, stream: {stream_status})") - + # CRITICAL FIX: Check for client disconnect during wait periods # Django/WSGI might not immediately detect disconnections, but we can check periodically if consecutive_empty > 10: # After some number of empty reads @@ -374,7 +399,7 @@ def stream_ts(request, channel_id): # Error reading from connection, likely closed logger.info(f"[{client_id}] Connection error, client likely disconnected") break - + # Disconnect after long inactivity # For non-owner workers, we're more lenient with timeout if time.time() - last_yield_time > Config.STREAM_TIMEOUT: @@ -385,25 +410,25 @@ def stream_ts(request, channel_id): # Non-owner worker without data for too long logger.warning(f"[{client_id}] Non-owner worker with no data for {Config.STREAM_TIMEOUT}s, disconnecting") break - + # ADD THIS: Check if worker has more recent chunks but still stuck # This can indicate the client is disconnected but we're not detecting it if consecutive_empty > 100 and buffer.index > local_index + 50: logger.warning(f"[{client_id}] Possible ghost client: buffer has advanced {buffer.index - local_index} chunks ahead but client stuck at {local_index}") break - + except Exception as e: logger.error(f"[{client_id}] Stream error: {e}", exc_info=True) finally: # Client cleanup elapsed = time.time() - stream_start_time local_clients = 0 - + if channel_id in proxy_server.client_managers: local_clients = proxy_server.client_managers[channel_id].remove_client(client_id) total_clients = proxy_server.client_managers[channel_id].get_total_client_count() logger.info(f"[{client_id}] Disconnected after {elapsed:.2f}s, {bytes_sent/1024:.1f}KB in {chunks_sent} chunks (local: {local_clients}, total: {total_clients})") - + # If no clients left and we're the owner, schedule shutdown using the config value if local_clients == 0 and proxy_server.am_i_owner(channel_id): logger.info(f"No local clients left for channel {channel_id}, scheduling shutdown") @@ -412,7 +437,7 @@ def stream_ts(request, channel_id): shutdown_delay = getattr(Config, 'CHANNEL_SHUTDOWN_DELAY', 5) logger.info(f"Waiting {shutdown_delay}s before checking if channel should be stopped") time.sleep(shutdown_delay) - + # After delay, check global client count if channel_id in proxy_server.client_managers: total = proxy_server.client_managers[channel_id].get_total_client_count() @@ -421,17 +446,17 @@ def stream_ts(request, channel_id): proxy_server.stop_channel(channel_id) else: logger.info(f"Not shutting down channel {channel_id}, {total} clients still connected") - + shutdown_thread = threading.Thread(target=delayed_shutdown) shutdown_thread.daemon = True shutdown_thread.start() - + response = StreamingHttpResponse( streaming_content=generate(), content_type='video/mp2t' ) return response - + except Exception as e: logger.error(f"Error in stream_ts: {e}", exc_info=True) return JsonResponse({'error': str(e)}, status=500) @@ -444,16 +469,16 @@ def change_stream(request, channel_id): data = json.loads(request.body) new_url = data.get('url') user_agent = data.get('user_agent') - + if not new_url: return JsonResponse({'error': 'No URL provided'}, status=400) - + logger.info(f"Attempting to change stream URL for channel {channel_id} to {new_url}") - + # Enhanced channel detection in_local_managers = channel_id in proxy_server.stream_managers in_local_buffers = channel_id in proxy_server.stream_buffers - + # First check Redis directly before using our wrapper method redis_keys = None if proxy_server.redis_client: @@ -462,17 +487,17 @@ def change_stream(request, channel_id): redis_keys = [k.decode('utf-8') for k in redis_keys] if redis_keys else [] except Exception as e: logger.error(f"Error checking Redis keys: {e}") - + # Now use our standard check channel_exists = proxy_server.check_if_channel_exists(channel_id) - + # Log detailed diagnostics logger.info(f"Channel {channel_id} diagnostics: " f"in_local_managers={in_local_managers}, " f"in_local_buffers={in_local_buffers}, " f"redis_keys_count={len(redis_keys) if redis_keys else 0}, " f"channel_exists={channel_exists}") - + if not channel_exists: # If channel doesn't exist but we found Redis keys, force initialize it if redis_keys: @@ -488,17 +513,17 @@ def change_stream(request, channel_id): 'redis_keys': redis_keys, } }, status=404) - + # Update metadata in Redis regardless of ownership - this ensures URL is updated # even if the owner worker is handling another request if proxy_server.redis_client: try: metadata_key = f"ts_proxy:channel:{channel_id}:metadata" - + # First check if the key exists and what type it is key_type = proxy_server.redis_client.type(metadata_key).decode('utf-8') logger.debug(f"Redis key {metadata_key} is of type: {key_type}") - + # Use the appropriate method based on the key type if key_type == 'hash': proxy_server.redis_client.hset(metadata_key, "url", new_url) @@ -517,21 +542,21 @@ def change_stream(request, channel_id): if user_agent: metadata["user_agent"] = user_agent proxy_server.redis_client.hset(metadata_key, mapping=metadata) - + # Set switch request flag to ensure all workers see it switch_key = f"ts_proxy:channel:{channel_id}:switch_request" proxy_server.redis_client.setex(switch_key, 30, new_url) # 30 second TTL - + logger.info(f"Updated metadata for channel {channel_id} in Redis") except Exception as e: logger.error(f"Error updating Redis metadata: {e}", exc_info=True) - + # If we're the owner, update directly if proxy_server.am_i_owner(channel_id) and channel_id in proxy_server.stream_managers: logger.info(f"This worker is the owner, changing stream URL for channel {channel_id}") manager = proxy_server.stream_managers[channel_id] old_url = manager.url - + # Update the stream result = manager.update_url(new_url) logger.info(f"Stream URL changed from {old_url} to {new_url}, result: {result}") @@ -542,7 +567,7 @@ def change_stream(request, channel_id): 'owner': True, 'worker_id': proxy_server.worker_id }) - + # If we're not the owner, publish an event for the owner to pick up else: logger.info(f"This worker is not the owner, requesting URL change via Redis PubSub") @@ -555,12 +580,12 @@ def change_stream(request, channel_id): "requester": proxy_server.worker_id, "timestamp": time.time() } - + proxy_server.redis_client.publish( - f"ts_proxy:events:{channel_id}", + f"ts_proxy:events:{channel_id}", json.dumps(switch_request) ) - + return JsonResponse({ 'message': 'Stream URL change requested', 'channel': channel_id, @@ -568,9 +593,9 @@ def change_stream(request, channel_id): 'owner': False, 'worker_id': proxy_server.worker_id }) - + except json.JSONDecodeError: return JsonResponse({'error': 'Invalid JSON'}, status=400) except Exception as e: logger.error(f"Failed to change stream: {e}", exc_info=True) - return JsonResponse({'error': str(e)}, status=500) \ No newline at end of file + return JsonResponse({'error': str(e)}, status=500) diff --git a/core/models.py b/core/models.py index 6f353126..2309fdf4 100644 --- a/core/models.py +++ b/core/models.py @@ -47,6 +47,14 @@ class StreamProfile(models.Model): def __str__(self): return self.profile_name + def build_command(self, stream_url): + cmd = [] + if self.command == "ffmpeg": + cmd = ["ffmpeg", "-i", stream_url] + self.parameters.split() + ["pipe:1"] + elif self.command == "streamlink": + cmd = ["streamlink", stream_url] + self.parameters.split() + + return cmd class CoreSettings(models.Model): key = models.CharField(