""" VOD (Video on Demand) proxy views for handling movie and series streaming. Supports M3U profiles for authentication and URL transformation. """ import time import random import logging import requests from django.http import JsonResponse, Http404, HttpResponse from django.shortcuts import get_object_or_404 from django.views.decorators.csrf import csrf_exempt from apps.vod.models import Movie, Series, Episode from apps.vod.catalog_settings import get_vod_catalog_account_id, get_vod_playback_allowed_account_ids, filter_m3u_relation_queryset_for_playback, filter_relation_qs_for_user_playback from apps.m3u.models import M3UAccountProfile from apps.proxy.vod_proxy.multi_worker_connection_manager import MultiWorkerVODConnectionManager, infer_content_type_from_url, get_vod_client_stop_key from .utils import get_client_info from rest_framework.decorators import api_view, permission_classes from rest_framework.permissions import AllowAny from apps.accounts.models import User from apps.accounts.permissions import IsAdmin from apps.proxy.utils import check_user_stream_limits from dispatcharr.utils import network_access_allowed from core.utils import RedisClient logger = logging.getLogger(__name__) def _playback_relations_queryset(content_obj, user=None): """ M3U relations eligible for streaming. Optional DISPATCHARR_VOD_PLAYBACK_ACCOUNT_IDS=10,11,12 limits failover to those accounts while catalog listings can stay single-account. Per-user `vod_playback_mode=strict` further restricts to the user's accounts. """ allowed = get_vod_playback_allowed_account_ids() qs = content_obj.m3u_relations.all() if allowed is not None: qs = qs.filter(m3u_account_id__in=allowed, m3u_account__is_active=True) else: qs = qs.filter(m3u_account__is_active=True) qs, _ = filter_relation_qs_for_user_playback(qs, user) return qs.select_related("m3u_account") _request_times = {} def _account_active_connections(account_ids): """ Return active connection counts per M3U account using shared profile counters. Counts include both live and VOD sessions. """ counts = {int(aid): 0 for aid in account_ids} if not account_ids: return counts try: redis_client = RedisClient.get_client() if not redis_client: return counts profiles = M3UAccountProfile.objects.filter( m3u_account_id__in=account_ids, is_active=True ).values("id", "m3u_account_id") for p in profiles: counts[int(p["m3u_account_id"])] += int( redis_client.get(f"profile_connections:{p['id']}") or 0 ) except Exception: return counts return counts def _preempt_oldest_vod_session_for_account(m3u_account_id: int) -> bool: """ Stop the oldest active VOD session for the given account. Returns True if a stop signal was sent. """ try: redis_client = RedisClient.get_client() if not redis_client: return False profile_ids = set( M3UAccountProfile.objects.filter( m3u_account_id=m3u_account_id, is_active=True ).values_list("id", flat=True) ) if not profile_ids: return False oldest_session_id = None oldest_created_at = None for key in redis_client.scan_iter(match="vod_persistent_connection:*", count=1000): key_str = key.decode() if isinstance(key, bytes) else str(key) session_id = key_str.split(":", 1)[1] if ":" in key_str else None if not session_id: continue data = redis_client.hgetall(key) if not data: continue profile_id = data.get("m3u_profile_id") try: profile_id = int(profile_id) if profile_id is not None else None except (TypeError, ValueError): profile_id = None if profile_id not in profile_ids: continue created_raw = data.get("created_at") try: created_at = float(created_raw) if created_raw is not None else time.time() except (TypeError, ValueError): created_at = time.time() if oldest_created_at is None or created_at < oldest_created_at: oldest_created_at = created_at oldest_session_id = session_id if not oldest_session_id: return False redis_client.setex(get_vod_client_stop_key(oldest_session_id), 60, "true") logger.warning( f"[VOD-PREEMPT] Stopping oldest VOD session {oldest_session_id} for m3u account {m3u_account_id}" ) return True except Exception as e: logger.warning(f"[VOD-PREEMPT] Could not preempt VOD session for account {m3u_account_id}: {e}") return False def _pick_relation_by_account_load(relations_query, prefer_m3u_account_id=None, user=None): """ Select relation by least-loaded account first, then user-preferred account, then catalog tie-break, then highest account priority. If user has playback prefs (mode=prefer/strict) the user's first account becomes the preferred one for tie-breaking; strict mode is already enforced upstream. """ relations = list(relations_query.select_related("m3u_account")) if not relations: return None user_prefs = filter_relation_qs_for_user_playback(relations_query.none(), user)[1] account_ids = {int(r.m3u_account_id) for r in relations if r.m3u_account_id} load_map = _account_active_connections(account_ids) pref_id = int(prefer_m3u_account_id) if prefer_m3u_account_id is not None else None user_pref_id = int(user_prefs) if user_prefs is not None else None def sort_key(r): aid = int(r.m3u_account_id or 0) load = load_map.get(aid, 0) user_rank = 0 if (user_pref_id is not None and aid == user_pref_id) else 1 pref_rank = 0 if (pref_id is not None and aid == pref_id) else 1 return (load, user_rank, pref_rank, -(r.m3u_account.priority or 0), r.id) return sorted(relations, key=sort_key)[0] def _relations_by_account_load(content_obj, prefer_m3u_account_id=None, user=None): """ Return all active relations for a content object ordered by account load, user-preferred account, catalog tie-break, then priority. """ try: relations = list(_playback_relations_queryset(content_obj, user=user)) except Exception: return [] if not relations: return [] user_prefs = filter_relation_qs_for_user_playback(content_obj.m3u_relations.none(), user)[1] user_pref_id = int(user_prefs) if user_prefs is not None else None account_ids = {int(r.m3u_account_id) for r in relations if r.m3u_account_id} load_map = _account_active_connections(account_ids) pref_id = int(prefer_m3u_account_id) if prefer_m3u_account_id is not None else None def sort_key(r): aid = int(r.m3u_account_id or 0) load = load_map.get(aid, 0) user_rank = 0 if (user_pref_id is not None and aid == user_pref_id) else 1 pref_rank = 0 if (pref_id is not None and aid == pref_id) else 1 return (load, user_rank, pref_rank, -(r.m3u_account.priority or 0), r.id) return sorted(relations, key=sort_key) def _get_content_and_relation(content_type, content_id, preferred_m3u_account_id=None, preferred_stream_id=None, user=None): """Get the content object and its M3U relation""" try: logger.info(f"[CONTENT-LOOKUP] Looking up {content_type} with UUID {content_id}") if preferred_m3u_account_id: logger.info(f"[CONTENT-LOOKUP] Preferred M3U account ID: {preferred_m3u_account_id}") if preferred_stream_id: logger.info(f"[CONTENT-LOOKUP] Preferred stream ID: {preferred_stream_id}") if content_type == 'movie': content_obj = get_object_or_404(Movie, uuid=content_id) logger.info(f"[CONTENT-FOUND] Movie: {content_obj.name} (ID: {content_obj.id})") # Filter by preferred stream ID first (most specific) relations_query = _playback_relations_queryset(content_obj, user=user) if preferred_stream_id: specific_relation = relations_query.filter(stream_id=preferred_stream_id).first() if specific_relation: logger.info(f"[STREAM-SELECTED] Using specific stream: {specific_relation.stream_id} from provider: {specific_relation.m3u_account.name}") return content_obj, specific_relation else: logger.warning(f"[STREAM-FALLBACK] Preferred stream ID {preferred_stream_id} not found, falling back to account/priority selection") # Filter by preferred M3U account if specified if preferred_m3u_account_id: specific_relation = relations_query.filter(m3u_account__id=preferred_m3u_account_id).first() if specific_relation: logger.info(f"[PROVIDER-SELECTED] Using preferred provider: {specific_relation.m3u_account.name}") return content_obj, specific_relation else: logger.warning(f"[PROVIDER-FALLBACK] Preferred M3U account {preferred_m3u_account_id} not found, using highest priority") # Fallback: choose least-loaded account (live + VOD), then priority relation = _pick_relation_by_account_load(relations_query, prefer_m3u_account_id=get_vod_catalog_account_id(), user=user) if relation: logger.info(f"[PROVIDER-SELECTED] Using provider: {relation.m3u_account.name} (priority: {relation.m3u_account.priority})") return content_obj, relation elif content_type == 'episode': content_obj = get_object_or_404(Episode, uuid=content_id) logger.info(f"[CONTENT-FOUND] Episode: {content_obj.name} (ID: {content_obj.id}, Series: {content_obj.series.name})") # Filter by preferred stream ID first (most specific) relations_query = _playback_relations_queryset(content_obj, user=user) if preferred_stream_id: specific_relation = relations_query.filter(stream_id=preferred_stream_id).first() if specific_relation: logger.info(f"[STREAM-SELECTED] Using specific stream: {specific_relation.stream_id} from provider: {specific_relation.m3u_account.name}") return content_obj, specific_relation else: logger.warning(f"[STREAM-FALLBACK] Preferred stream ID {preferred_stream_id} not found, falling back to account/priority selection") # Filter by preferred M3U account if specified if preferred_m3u_account_id: specific_relation = relations_query.filter(m3u_account__id=preferred_m3u_account_id).first() if specific_relation: logger.info(f"[PROVIDER-SELECTED] Using preferred provider: {specific_relation.m3u_account.name}") return content_obj, specific_relation else: logger.warning(f"[PROVIDER-FALLBACK] Preferred M3U account {preferred_m3u_account_id} not found, using highest priority") # Fallback: choose least-loaded account (live + VOD), then priority relation = _pick_relation_by_account_load(relations_query, prefer_m3u_account_id=get_vod_catalog_account_id(), user=user) if relation: logger.info(f"[PROVIDER-SELECTED] Using provider: {relation.m3u_account.name} (priority: {relation.m3u_account.priority})") return content_obj, relation elif content_type == 'series': # For series, get the first episode series = get_object_or_404(Series, uuid=content_id) logger.info(f"[CONTENT-FOUND] Series: {series.name} (ID: {series.id})") episode = series.episodes.first() if not episode: logger.error(f"[CONTENT-ERROR] No episodes found for series {series.name}") return None, None logger.info(f"[CONTENT-FOUND] First episode: {episode.name} (ID: {episode.id})") # Filter by preferred stream ID first (most specific) relations_query = _playback_relations_queryset(episode, user=user) if preferred_stream_id: specific_relation = relations_query.filter(stream_id=preferred_stream_id).first() if specific_relation: logger.info(f"[STREAM-SELECTED] Using specific stream: {specific_relation.stream_id} from provider: {specific_relation.m3u_account.name}") return episode, specific_relation else: logger.warning(f"[STREAM-FALLBACK] Preferred stream ID {preferred_stream_id} not found, falling back to account/priority selection") # Filter by preferred M3U account if specified if preferred_m3u_account_id: specific_relation = relations_query.filter(m3u_account__id=preferred_m3u_account_id).first() if specific_relation: logger.info(f"[PROVIDER-SELECTED] Using preferred provider: {specific_relation.m3u_account.name}") return episode, specific_relation else: logger.warning(f"[PROVIDER-FALLBACK] Preferred M3U account {preferred_m3u_account_id} not found, using highest priority") # Fallback: choose least-loaded account (live + VOD), then priority relation = _pick_relation_by_account_load(relations_query, prefer_m3u_account_id=get_vod_catalog_account_id(), user=user) if relation: logger.info(f"[PROVIDER-SELECTED] Using provider: {relation.m3u_account.name} (priority: {relation.m3u_account.priority})") return episode, relation else: logger.error(f"[CONTENT-ERROR] Invalid content type: {content_type}") return None, None except Exception as e: logger.error(f"Error getting content object: {e}") return None, None def _get_stream_url_from_relation(relation): """Get stream URL from the M3U relation""" try: # Log the relation type and available attributes logger.info(f"[VOD-URL] Relation type: {type(relation).__name__}") logger.info(f"[VOD-URL] Account type: {relation.m3u_account.account_type}") logger.info(f"[VOD-URL] Stream ID: {getattr(relation, 'stream_id', 'N/A')}") # First try the get_stream_url method (this should build URLs dynamically) if hasattr(relation, 'get_stream_url'): url = relation.get_stream_url() if url: logger.info(f"[VOD-URL] Built URL from get_stream_url(): {url}") return url else: logger.warning(f"[VOD-URL] get_stream_url() returned None") logger.error(f"[VOD-URL] Relation has no get_stream_url method or it failed") return None except Exception as e: logger.error(f"[VOD-URL] Error getting stream URL from relation: {e}", exc_info=True) return None def _get_m3u_profile(m3u_account, profile_id, session_id=None): """Get appropriate M3U profile for streaming using Redis-based viewer counts Args: m3u_account: M3UAccount instance profile_id: Optional specific profile ID requested session_id: Optional session ID to check for existing connections Returns: tuple: (M3UAccountProfile, current_connections) or None if no profile found """ try: from core.utils import RedisClient redis_client = RedisClient.get_client() if not redis_client: logger.warning("Redis not available, falling back to default profile") default_profile = M3UAccountProfile.objects.filter( m3u_account=m3u_account, is_active=True, is_default=True ).first() return (default_profile, 0) if default_profile else None # Check if this session already has an active connection if session_id: persistent_connection_key = f"vod_persistent_connection:{session_id}" connection_data = redis_client.hgetall(persistent_connection_key) if connection_data: existing_profile_id = connection_data.get('m3u_profile_id') if existing_profile_id: try: existing_profile = M3UAccountProfile.objects.get( id=int(existing_profile_id), m3u_account=m3u_account, is_active=True ) # Get current connections for logging profile_connections_key = f"profile_connections:{existing_profile.id}" current_connections = int(redis_client.get(profile_connections_key) or 0) logger.info(f"[PROFILE-SELECTION] Session {session_id} reusing existing profile {existing_profile.id}: {current_connections}/{existing_profile.max_streams} connections") return (existing_profile, current_connections) except (M3UAccountProfile.DoesNotExist, ValueError): logger.warning(f"[PROFILE-SELECTION] Session {session_id} has invalid profile ID {existing_profile_id}, selecting new profile") except Exception as e: logger.warning(f"[PROFILE-SELECTION] Error checking existing profile for session {session_id}: {e}") else: logger.debug(f"[PROFILE-SELECTION] Session {session_id} exists but has no profile ID stored") # If specific profile requested, try to use it if profile_id: try: profile = M3UAccountProfile.objects.get( id=profile_id, m3u_account=m3u_account, is_active=True ) # Check Redis-based current connections profile_connections_key = f"profile_connections:{profile.id}" current_connections = int(redis_client.get(profile_connections_key) or 0) if profile.max_streams == 0 or current_connections < profile.max_streams: logger.info(f"[PROFILE-SELECTION] Using requested profile {profile.id}: {current_connections}/{profile.max_streams} connections") return (profile, current_connections) else: logger.warning(f"[PROFILE-SELECTION] Requested profile {profile.id} is at capacity: {current_connections}/{profile.max_streams}") except M3UAccountProfile.DoesNotExist: logger.warning(f"[PROFILE-SELECTION] Requested profile {profile_id} not found") # Get active profiles ordered by priority (default first) m3u_profiles = M3UAccountProfile.objects.filter( m3u_account=m3u_account, is_active=True ) default_profile = m3u_profiles.filter(is_default=True).first() if not default_profile: logger.error(f"[PROFILE-SELECTION] No default profile found for M3U account {m3u_account.id}") return None # Check profiles in order: default first, then others profiles = [default_profile] + list(m3u_profiles.filter(is_default=False)) for profile in profiles: profile_connections_key = f"profile_connections:{profile.id}" current_connections = int(redis_client.get(profile_connections_key) or 0) # Check if profile has available connection slots if profile.max_streams == 0 or current_connections < profile.max_streams: logger.info(f"[PROFILE-SELECTION] Selected profile {profile.id} ({profile.name}): {current_connections}/{profile.max_streams} connections") return (profile, current_connections) else: logger.debug(f"[PROFILE-SELECTION] Profile {profile.id} at capacity: {current_connections}/{profile.max_streams}") # All profiles are at capacity: preempt oldest VOD session on this account once, then retry logger.error(f"[PROFILE-SELECTION] All profiles at capacity for M3U account {m3u_account.id}, rejecting request") if _preempt_oldest_vod_session_for_account(int(m3u_account.id)): # Give workers a brief chance to release counters time.sleep(0.2) for profile in profiles: profile_connections_key = f"profile_connections:{profile.id}" current_connections = int(redis_client.get(profile_connections_key) or 0) if profile.max_streams == 0 or current_connections < profile.max_streams: logger.info( f"[PROFILE-SELECTION] Recovered slot after preemption with profile {profile.id}: {current_connections}/{profile.max_streams}" ) return (profile, current_connections) return None except Exception as e: logger.error(f"Error getting M3U profile: {e}") return None def _transform_url(original_url, m3u_profile): """Transform URL based on M3U profile settings""" try: import regex if not original_url: return None search_pattern = m3u_profile.search_pattern replace_pattern = m3u_profile.replace_pattern # Convert JS-style backreferences in replace: $ -> \g, $1 -> \1 safe_replace_pattern = regex.sub(r'\$<([^>]+)>', r'\\g<\1>', replace_pattern) safe_replace_pattern = regex.sub(r'\$(\d+)', r'\\\1', safe_replace_pattern) if search_pattern and replace_pattern: # regex module accepts JS-style (?...) named groups natively transformed_url = regex.sub(search_pattern, safe_replace_pattern, original_url) return transformed_url return original_url except Exception as e: logger.error(f"Error transforming URL: {e}") return original_url @api_view(["GET"]) @permission_classes([AllowAny]) def stream_vod(request, content_type, content_id, session_id=None, profile_id=None, user=None): """ Stream VOD content (movies or series episodes) with session-based connection reuse Args: content_type: 'movie', 'series', or 'episode' content_id: ID of the content session_id: Optional session ID from URL path (for persistent connections) profile_id: Optional M3U profile ID for authentication """ if not network_access_allowed(request, "STREAMS"): return JsonResponse({"error": "Forbidden"}, status=403) logger.info(f"[VOD-REQUEST] Starting VOD stream request: {content_type}/{content_id}, session: {session_id}, profile: {profile_id}") logger.info(f"[VOD-REQUEST] Full request path: {request.get_full_path()}") logger.info(f"[VOD-REQUEST] Request method: {request.method}") logger.info(f"[VOD-REQUEST] Request headers: {dict(request.headers)}") try: if user is None: req_user = getattr(request, "user", None) if req_user is not None and getattr(req_user, "is_authenticated", False): user = req_user else: qs_user_id = request.GET.get("user_id") if qs_user_id and str(qs_user_id).isdigit(): try: user = User.objects.filter(id=int(qs_user_id)).first() except Exception: user = None client_ip, client_user_agent = get_client_info(request) # Extract timeshift parameters from query string # Support multiple timeshift parameter formats utc_start = request.GET.get('utc_start') or request.GET.get('start') or request.GET.get('playliststart') utc_end = request.GET.get('utc_end') or request.GET.get('end') or request.GET.get('playlistend') offset = request.GET.get('offset') or request.GET.get('seek') or request.GET.get('t') # VLC specific timeshift parameters if not utc_start and not offset: # Check for VLC-style timestamp parameters if 'timestamp' in request.GET: offset = request.GET.get('timestamp') elif 'time' in request.GET: offset = request.GET.get('time') # Session ID now comes from URL path parameter # Remove legacy query parameter extraction since we're using path-based routing # Extract Range header for seeking support range_header = request.META.get('HTTP_RANGE') logger.info(f"[VOD-TIMESHIFT] Timeshift params - utc_start: {utc_start}, utc_end: {utc_end}, offset: {offset}") logger.info(f"[VOD-SESSION] Session ID: {session_id}") # Log all query parameters for debugging if request.GET: logger.debug(f"[VOD-PARAMS] All query params: {dict(request.GET)}") if range_header: logger.info(f"[VOD-RANGE] Range header: {range_header}") # Parse the range to understand what position VLC is seeking to try: if 'bytes=' in range_header: range_part = range_header.replace('bytes=', '') if '-' in range_part: start_byte, end_byte = range_part.split('-', 1) if start_byte: start_pos_mb = int(start_byte) / (1024 * 1024) logger.info(f"[VOD-SEEK] Seeking to byte position: {start_byte} (~{start_pos_mb:.1f} MB)") if int(start_byte) > 0: logger.info(f"[VOD-SEEK] *** ACTUAL SEEK DETECTED *** Position: {start_pos_mb:.1f} MB") else: logger.info(f"[VOD-SEEK] Open-en`ded range request (from start)") if end_byte: end_pos_mb = int(end_byte) / (1024 * 1024) logger.info(f"[VOD-SEEK] End position: {end_byte} bytes (~{end_pos_mb:.1f} MB)") except Exception as e: logger.warning(f"[VOD-SEEK] Could not parse range header: {e}") # Simple seek detection - track rapid requests current_time = time.time() request_key = f"{client_ip}:{content_type}:{content_id}" if request_key in _request_times: time_diff = current_time - _request_times[request_key] if time_diff < 5.0: logger.info(f"[VOD-SEEK] Rapid request detected ({time_diff:.1f}s) - likely seeking") _request_times[request_key] = current_time else: logger.info(f"[VOD-RANGE] No Range header - full content request") logger.info(f"[VOD-CLIENT] Client info - IP: {client_ip}, User-Agent: {client_user_agent[:50]}...") # If no session ID, create one and redirect to path-based URL if not session_id: new_session_id = f"vod_{int(time.time() * 1000)}_{random.randint(1000, 9999)}" logger.info(f"[VOD-SESSION] Creating new session: {new_session_id}") # Preserve any query parameters (except session_id) query_params = dict(request.GET) query_params.pop('session_id', None) # Remove if present if user: redirect_url = f"{request.path}?session_id={new_session_id}" if query_params: query_string = urlencode(query_params, doseq=True) redirect_url = f"{redirect_url}&{query_string}" else: # Build redirect URL with session ID in path, preserve query parameters path_parts = request.path.rstrip('/').split('/') # Construct new path: /vod/movie/UUID/SESSION_ID or /vod/movie/UUID/SESSION_ID/PROFILE_ID/ if profile_id: new_path = f"{'/'.join(path_parts)}/{new_session_id}/{profile_id}/" else: new_path = f"{'/'.join(path_parts)}/{new_session_id}" if query_params: from urllib.parse import urlencode query_string = urlencode(query_params, doseq=True) redirect_url = f"{new_path}?{query_string}" else: redirect_url = new_path logger.info(f"[VOD-SESSION] Redirecting to path-based URL: {redirect_url}") return HttpResponse( status=301, headers={'Location': redirect_url} ) if user: if not check_user_stream_limits(user, session_id, media_id=content_id): return JsonResponse( {"error": f"Stream limit exceeded ({user.stream_limit} concurrent streams allowed)"}, status=429 ) # Extract preferred M3U account ID and stream ID from query parameters preferred_m3u_account_id = request.GET.get('m3u_account_id') preferred_stream_id = request.GET.get('stream_id') if preferred_m3u_account_id: try: preferred_m3u_account_id = int(preferred_m3u_account_id) except (ValueError, TypeError): logger.warning(f"[VOD-PARAM] Invalid m3u_account_id parameter: {preferred_m3u_account_id}") preferred_m3u_account_id = None if preferred_stream_id: logger.info(f"[VOD-PARAM] Preferred stream ID: {preferred_stream_id}") # Get the content object and its relation content_obj, relation = _get_content_and_relation(content_type, content_id, preferred_m3u_account_id, preferred_stream_id, user=user) if not content_obj or not relation: logger.error(f"[VOD-ERROR] Content or relation not found: {content_type} {content_id}") raise Http404(f"Content not found: {content_type} {content_id}") logger.info(f"[VOD-CONTENT] Found content: {getattr(content_obj, 'name', 'Unknown')}") # Get M3U account from relation m3u_account = relation.m3u_account logger.info(f"[VOD-ACCOUNT] Using M3U account: {m3u_account.name}") # Get stream URL from relation stream_url = _get_stream_url_from_relation(relation) logger.info(f"[VOD-CONTENT] Content URL: {stream_url or 'No URL found'}") if not stream_url: logger.error(f"[VOD-ERROR] No stream URL available for {content_type} {content_id}") return HttpResponse("No stream URL available", status=503) # Get M3U profile (returns profile and current connection count) profile_result = _get_m3u_profile(m3u_account, profile_id, session_id) # Auto-mode fallback: if selected account is at capacity, try next accounts if (not profile_result or not profile_result[0]) and not preferred_stream_id: logger.warning( f"[VOD-FALLBACK] Initial account {m3u_account.name} at capacity, trying alternate accounts" ) for alt_relation in _relations_by_account_load(content_obj, prefer_m3u_account_id=get_vod_catalog_account_id(), user=user): if alt_relation.id == relation.id: continue alt_account = alt_relation.m3u_account alt_profile_result = _get_m3u_profile(alt_account, profile_id, session_id) if not alt_profile_result or not alt_profile_result[0]: continue alt_stream_url = _get_stream_url_from_relation(alt_relation) if not alt_stream_url: continue logger.info( f"[VOD-FALLBACK] Switched to account {alt_account.name} for {content_type} {content_id}" ) relation = alt_relation m3u_account = alt_account stream_url = alt_stream_url profile_result = alt_profile_result break if not profile_result or not profile_result[0]: logger.error(f"[VOD-ERROR] No suitable M3U profile found for {content_type} {content_id}") return HttpResponse("No available stream", status=503) m3u_profile, current_connections = profile_result logger.info(f"[VOD-PROFILE] Using M3U profile: {m3u_profile.id} (max_streams: {m3u_profile.max_streams}, current: {current_connections})") # Connection tracking is handled by the connection manager # Transform URL based on profile final_stream_url = _transform_url(stream_url, m3u_profile) logger.info(f"[VOD-URL] Final stream URL: {final_stream_url}") # Validate stream URL if not final_stream_url or not final_stream_url.startswith(('http://', 'https://')): logger.error(f"[VOD-ERROR] Invalid stream URL: {final_stream_url}") return HttpResponse("Invalid stream URL", status=500) # Get connection manager (Redis-backed for multi-worker support) connection_manager = MultiWorkerVODConnectionManager.get_instance() # Stream the content with session-based connection reuse logger.info("[VOD-STREAM] Calling connection manager to stream content") response = connection_manager.stream_content_with_session( session_id=session_id, content_obj=content_obj, stream_url=final_stream_url, m3u_profile=m3u_profile, client_ip=client_ip, client_user_agent=client_user_agent, request=request, utc_start=utc_start, utc_end=utc_end, offset=offset, range_header=range_header, user=user, ) logger.info(f"[VOD-SUCCESS] Stream response created successfully, type: {type(response)}") return response except Exception as e: logger.error(f"[VOD-EXCEPTION] Error streaming {content_type} {content_id}: {e}", exc_info=True) return HttpResponse(f"Streaming error: {str(e)}", status=500) @api_view(["HEAD"]) @permission_classes([AllowAny]) def head_vod(request, content_type, content_id, session_id=None, profile_id=None): """ Handle HEAD requests for FUSE filesystem integration Returns content length and session URL header for subsequent GET requests """ if not network_access_allowed(request, "STREAMS"): return JsonResponse({"error": "Forbidden"}, status=403) logger.info(f"[VOD-HEAD] HEAD request: {content_type}/{content_id}, session: {session_id}, profile: {profile_id}") try: user = request.user if getattr(request, 'user', None) and getattr(request.user, 'is_authenticated', False) else None # Get client info for M3U profile selection client_ip, client_user_agent = get_client_info(request) logger.info(f"[VOD-HEAD] Client info - IP: {client_ip}, User-Agent: {client_user_agent[:50] if client_user_agent else 'None'}...") # If no session ID, create one (same logic as GET) if not session_id: new_session_id = f"vod_{int(time.time() * 1000)}_{random.randint(1000, 9999)}" logger.info(f"[VOD-HEAD] Creating new session for HEAD: {new_session_id}") # Build session URL for response header path_parts = request.path.rstrip('/').split('/') if profile_id: session_url = f"{'/'.join(path_parts)}/{new_session_id}/{profile_id}/" else: session_url = f"{'/'.join(path_parts)}/{new_session_id}" session_id = new_session_id else: # Session already in URL, construct the current session URL session_url = request.path logger.info(f"[VOD-HEAD] Using existing session: {session_id}") # Extract preferred M3U account ID and stream ID from query parameters preferred_m3u_account_id = request.GET.get('m3u_account_id') preferred_stream_id = request.GET.get('stream_id') if preferred_m3u_account_id: try: preferred_m3u_account_id = int(preferred_m3u_account_id) except (ValueError, TypeError): logger.warning(f"[VOD-HEAD] Invalid m3u_account_id parameter: {preferred_m3u_account_id}") preferred_m3u_account_id = None if preferred_stream_id: logger.info(f"[VOD-HEAD] Preferred stream ID: {preferred_stream_id}") # Get content and relation (same as GET) content_obj, relation = _get_content_and_relation(content_type, content_id, preferred_m3u_account_id, preferred_stream_id, user=user) if not content_obj or not relation: logger.error(f"[VOD-HEAD] Content or relation not found: {content_type} {content_id}") return HttpResponse("Content not found", status=404) # Get M3U account and stream URL m3u_account = relation.m3u_account stream_url = _get_stream_url_from_relation(relation) if not stream_url: logger.error(f"[VOD-HEAD] No stream URL available for {content_type} {content_id}") return HttpResponse("No stream URL available", status=503) # Get M3U profile (returns profile and current connection count) profile_result = _get_m3u_profile(m3u_account, profile_id, session_id) if (not profile_result or not profile_result[0]) and not preferred_stream_id: logger.warning( f"[VOD-HEAD] Initial account {m3u_account.name} at capacity, trying alternate accounts" ) for alt_relation in _relations_by_account_load(content_obj, prefer_m3u_account_id=get_vod_catalog_account_id(), user=user): if alt_relation.id == relation.id: continue alt_account = alt_relation.m3u_account alt_profile_result = _get_m3u_profile(alt_account, profile_id, session_id) if not alt_profile_result or not alt_profile_result[0]: continue alt_stream_url = _get_stream_url_from_relation(alt_relation) if not alt_stream_url: continue relation = alt_relation m3u_account = alt_account stream_url = alt_stream_url profile_result = alt_profile_result logger.info(f"[VOD-HEAD] Fallback switched to account {alt_account.name}") break if not profile_result or not profile_result[0]: logger.error(f"[VOD-HEAD] No M3U profile found or all profiles at capacity") return HttpResponse("No available stream", status=503) m3u_profile, current_connections = profile_result # Transform URL if needed final_stream_url = _transform_url(stream_url, m3u_profile) # Make a small range GET request to get content length since providers don't support HEAD # We'll use a tiny range to minimize data transfer but get the headers we need # Use M3U account's user agent as primary, client user agent as fallback m3u_user_agent = m3u_account.get_user_agent().user_agent if m3u_account.get_user_agent() else None headers = { 'User-Agent': m3u_user_agent or client_user_agent or 'Dispatcharr/1.0', 'Accept': '*/*', 'Range': 'bytes=0-1' # Request only first 2 bytes } logger.info(f"[VOD-HEAD] Making small range GET request to provider: {final_stream_url}") response = requests.get(final_stream_url, headers=headers, timeout=30, allow_redirects=True, stream=True) # Check for range support - should be 206 for partial content if response.status_code == 206: # Parse Content-Range header to get total file size content_range = response.headers.get('Content-Range', '') if content_range: # Content-Range: bytes 0-1/1234567890 total_size = content_range.split('/')[-1] logger.info(f"[VOD-HEAD] Got file size from Content-Range: {total_size}") else: logger.warning(f"[VOD-HEAD] No Content-Range header in 206 response") total_size = response.headers.get('Content-Length', '0') elif response.status_code == 200: # Server doesn't support range requests, use Content-Length from full response total_size = response.headers.get('Content-Length', '0') logger.info(f"[VOD-HEAD] Server doesn't support ranges, got Content-Length: {total_size}") else: logger.error(f"[VOD-HEAD] Provider GET request failed: {response.status_code}") return HttpResponse("Provider error", status=response.status_code) # Close the small range request - we don't need to keep this connection response.close() # Store the total content length in Redis for the persistent connection to use try: import redis from django.conf import settings redis_host = getattr(settings, 'REDIS_HOST', 'localhost') redis_port = int(getattr(settings, 'REDIS_PORT', 6379)) redis_db = int(getattr(settings, 'REDIS_DB', 0)) redis_password = getattr(settings, 'REDIS_PASSWORD', '') redis_user = getattr(settings, 'REDIS_USER', '') ssl_params = getattr(settings, 'REDIS_SSL_PARAMS', {}) r = redis.StrictRedis( host=redis_host, port=redis_port, db=redis_db, password=redis_password if redis_password else None, username=redis_user if redis_user else None, decode_responses=True, **ssl_params ) content_length_key = f"vod_content_length:{session_id}" r.set(content_length_key, total_size, ex=1800) # Store for 30 minutes logger.info(f"[VOD-HEAD] Stored total content length {total_size} for session {session_id}") except Exception as e: logger.error(f"[VOD-HEAD] Failed to store content length in Redis: {e}") # Now create a persistent connection for the session (if one doesn't exist) # This ensures the FUSE GET requests will reuse the same connection connection_manager = MultiWorkerVODConnectionManager.get_instance() logger.info(f"[VOD-HEAD] Pre-creating persistent connection for session: {session_id}") # We don't actually stream content here, just ensure connection is ready # The actual GET requests from FUSE will use the persistent connection # Use the total_size we extracted from the range response provider_content_type = response.headers.get('Content-Type') if provider_content_type: content_type_header = provider_content_type logger.info(f"[VOD-HEAD] Using provider Content-Type: {content_type_header}") else: # Provider didn't send Content-Type, infer from URL inferred_content_type = infer_content_type_from_url(final_stream_url) if inferred_content_type: content_type_header = inferred_content_type logger.info(f"[VOD-HEAD] Provider missing Content-Type, inferred from URL: {content_type_header}") else: content_type_header = 'video/mp4' logger.info(f"[VOD-HEAD] No Content-Type from provider and could not infer from URL, using default: {content_type_header}") logger.info(f"[VOD-HEAD] Provider response - Total Size: {total_size}, Type: {content_type_header}") # Create response with content length and session URL header head_response = HttpResponse() head_response['Content-Length'] = total_size head_response['Content-Type'] = content_type_header head_response['Accept-Ranges'] = 'bytes' # Custom header with session URL for FUSE head_response['X-Session-URL'] = session_url head_response['X-Dispatcharr-Session'] = session_id logger.info(f"[VOD-HEAD] Returning HEAD response with session URL: {session_url}") return head_response except Exception as e: logger.error(f"[VOD-HEAD] Error in HEAD request: {e}", exc_info=True) return HttpResponse(f"HEAD error: {str(e)}", status=500) @api_view(["GET"]) @permission_classes([IsAdmin]) def vod_stats(request): """Get current VOD connection statistics""" try: connection_manager = MultiWorkerVODConnectionManager.get_instance() redis_client = connection_manager.redis_client if not redis_client: return JsonResponse({'error': 'Redis not available'}, status=500) # Get all VOD persistent connections (consolidated data) pattern = "vod_persistent_connection:*" cursor = 0 connections = [] current_time = time.time() while True: cursor, keys = redis_client.scan(cursor, match=pattern, count=100) for key in keys: try: connection_data = redis_client.hgetall(key) if connection_data: # Extract session ID from key session_id = key.replace('vod_persistent_connection:', '') # Decode Redis hash data combined_data = {} for k, v in connection_data.items(): combined_data[k] = v # Get content info from the connection data (using correct field names) content_type = combined_data.get('content_obj_type', 'unknown') content_uuid = combined_data.get('content_uuid', 'unknown') client_id = session_id # Get content info with enhanced metadata content_name = "Unknown" content_metadata = {} try: if content_type == 'movie': content_obj = Movie.objects.select_related('logo').get(uuid=content_uuid) content_name = content_obj.name # Get duration from content object duration_secs = None if hasattr(content_obj, 'duration_secs') and content_obj.duration_secs: duration_secs = content_obj.duration_secs # If we don't have duration_secs, try to calculate it from file size and position data if not duration_secs: file_size_bytes = int(combined_data.get('total_content_size', 0)) last_seek_byte = int(combined_data.get('last_seek_byte', 0)) last_seek_percentage = float(combined_data.get('last_seek_percentage', 0.0)) # Calculate position if we have the required data if file_size_bytes and file_size_bytes > 0 and last_seek_percentage > 0: # If we know the seek percentage and current time position, we can estimate duration # But we need to know the current time position in seconds first # For now, let's use a rough estimate based on file size and typical bitrates # This is a fallback - ideally duration should be in the database estimated_duration = 6000 # 100 minutes as default for movies duration_secs = estimated_duration content_metadata = { 'year': content_obj.year, 'rating': content_obj.rating, 'genre': content_obj.genre, 'duration_secs': duration_secs, 'description': content_obj.description, 'logo_url': content_obj.logo.url if content_obj.logo else None, 'tmdb_id': content_obj.tmdb_id, 'imdb_id': content_obj.imdb_id } elif content_type == 'episode': content_obj = Episode.objects.select_related('series', 'series__logo').get(uuid=content_uuid) content_name = f"{content_obj.series.name} - {content_obj.name}" # Get duration from content object duration_secs = None if hasattr(content_obj, 'duration_secs') and content_obj.duration_secs: duration_secs = content_obj.duration_secs # If we don't have duration_secs, estimate for episodes if not duration_secs: estimated_duration = 2400 # 40 minutes as default for episodes duration_secs = estimated_duration content_metadata = { 'series_name': content_obj.series.name, 'episode_name': content_obj.name, 'season_number': content_obj.season_number, 'episode_number': content_obj.episode_number, 'air_date': content_obj.air_date.isoformat() if content_obj.air_date else None, 'rating': content_obj.rating, 'duration_secs': duration_secs, 'description': content_obj.description, 'logo_url': content_obj.series.logo.url if content_obj.series.logo else None, 'series_year': content_obj.series.year, 'series_genre': content_obj.series.genre, 'tmdb_id': content_obj.tmdb_id, 'imdb_id': content_obj.imdb_id } except: pass # Get M3U profile information m3u_profile_info = {} m3u_profile_id = combined_data.get('m3u_profile_id') if m3u_profile_id: try: from apps.m3u.models import M3UAccountProfile profile = M3UAccountProfile.objects.select_related('m3u_account').get(id=m3u_profile_id) m3u_profile_info = { 'profile_name': profile.name, 'account_name': profile.m3u_account.name, 'account_id': profile.m3u_account.id, 'max_streams': profile.m3u_account.max_streams, 'm3u_profile_id': int(m3u_profile_id) } except Exception as e: logger.warning(f"Could not fetch M3U profile {m3u_profile_id}: {e}") # Also try to get profile info from stored data if database lookup fails if not m3u_profile_info and combined_data.get('m3u_profile_name'): m3u_profile_info = { 'profile_name': combined_data.get('m3u_profile_name', 'Unknown Profile'), 'm3u_profile_id': combined_data.get('m3u_profile_id'), 'account_name': 'Unknown Account' # We don't store account name directly } # Calculate estimated current position based on seek percentage or last known position last_known_position = int(combined_data.get('position_seconds', 0)) last_position_update = combined_data.get('last_position_update') last_seek_percentage = float(combined_data.get('last_seek_percentage', 0.0)) last_seek_timestamp = float(combined_data.get('last_seek_timestamp', 0.0)) estimated_position = last_known_position # If we have seek percentage and content duration, calculate position from that if last_seek_percentage > 0 and content_metadata.get('duration_secs'): try: duration_secs = int(content_metadata['duration_secs']) # Calculate position from seek percentage seek_position = int((last_seek_percentage / 100) * duration_secs) # If we have a recent seek timestamp, add elapsed time since seek if last_seek_timestamp > 0: elapsed_since_seek = current_time - last_seek_timestamp # Add elapsed time but don't exceed content duration estimated_position = min( seek_position + int(elapsed_since_seek), duration_secs ) else: estimated_position = seek_position except (ValueError, TypeError): pass elif last_position_update and content_metadata.get('duration_secs'): # Fallback: use time-based estimation from position_seconds try: update_timestamp = float(last_position_update) elapsed_since_update = current_time - update_timestamp # Add elapsed time to last known position, but don't exceed content duration estimated_position = min( last_known_position + int(elapsed_since_update), int(content_metadata['duration_secs']) ) except (ValueError, TypeError): # If timestamp parsing fails, fall back to last known position estimated_position = last_known_position connection_info = { 'content_type': content_type, 'content_uuid': content_uuid, 'content_name': content_name, 'content_metadata': content_metadata, 'm3u_profile': m3u_profile_info, 'client_id': client_id, 'client_ip': combined_data.get('client_ip', 'Unknown'), 'user_id': combined_data.get('user_id', '0'), 'user_agent': combined_data.get('client_user_agent', 'Unknown'), 'connected_at': combined_data.get('created_at'), 'last_activity': combined_data.get('last_activity'), 'm3u_profile_id': m3u_profile_id, 'position_seconds': estimated_position, # Use estimated position 'last_known_position': last_known_position, # Include raw position for debugging 'last_position_update': last_position_update, # Include timestamp for frontend use 'bytes_sent': int(combined_data.get('bytes_sent', 0)), # Seek/range information for position calculation and frontend display 'last_seek_byte': int(combined_data.get('last_seek_byte', 0)), 'last_seek_percentage': float(combined_data.get('last_seek_percentage', 0.0)), 'total_content_size': int(combined_data.get('total_content_size', 0)), 'last_seek_timestamp': float(combined_data.get('last_seek_timestamp', 0.0)) } # Calculate connection duration duration_calculated = False if connection_info['connected_at']: try: connected_time = float(connection_info['connected_at']) duration = current_time - connected_time connection_info['duration'] = int(duration) duration_calculated = True except: pass # Fallback: use last_activity if connected_at is not available if not duration_calculated and connection_info['last_activity']: try: last_activity_time = float(connection_info['last_activity']) # Estimate connection duration using client_id timestamp if available if connection_info['client_id'].startswith('vod_'): # Extract timestamp from client_id (format: vod_timestamp_random) parts = connection_info['client_id'].split('_') if len(parts) >= 2: client_start_time = float(parts[1]) / 1000.0 # Convert ms to seconds duration = current_time - client_start_time connection_info['duration'] = int(duration) duration_calculated = True except: pass # Final fallback if not duration_calculated: connection_info['duration'] = 0 connections.append(connection_info) except Exception as e: logger.error(f"Error processing connection key {key}: {e}") if cursor == 0: break # Group connections by content content_stats = {} for conn in connections: content_key = f"{conn['content_type']}:{conn['content_uuid']}" if content_key not in content_stats: content_stats[content_key] = { 'content_type': conn['content_type'], 'content_name': conn['content_name'], 'content_uuid': conn['content_uuid'], 'content_metadata': conn['content_metadata'], 'connection_count': 0, 'connections': [] } content_stats[content_key]['connection_count'] += 1 content_stats[content_key]['connections'].append(conn) return JsonResponse({ 'vod_connections': list(content_stats.values()), 'total_connections': len(connections), 'timestamp': current_time }) except Exception as e: logger.error(f"Error getting VOD stats: {e}") return JsonResponse({'error': str(e)}, status=500) @csrf_exempt @api_view(["POST"]) @permission_classes([IsAdmin]) def stop_vod_client(request): """Stop a specific VOD client connection using stop signal mechanism""" try: # Parse request body import json try: data = json.loads(request.body) except json.JSONDecodeError: return JsonResponse({'error': 'Invalid JSON'}, status=400) client_id = data.get('client_id') if not client_id: return JsonResponse({'error': 'No client_id provided'}, status=400) logger.info(f"Request to stop VOD client: {client_id}") # Get Redis client connection_manager = MultiWorkerVODConnectionManager.get_instance() redis_client = connection_manager.redis_client if not redis_client: return JsonResponse({'error': 'Redis not available'}, status=500) # Check if connection exists connection_key = f"vod_persistent_connection:{client_id}" connection_data = redis_client.hgetall(connection_key) if not connection_data: logger.warning(f"VOD connection not found: {client_id}") return JsonResponse({'error': 'Connection not found'}, status=404) # Set a stop signal key that the worker will check stop_key = get_vod_client_stop_key(client_id) redis_client.setex(stop_key, 60, "true") # 60 second TTL logger.info(f"Set stop signal for VOD client: {client_id}") return JsonResponse({ 'message': 'VOD client stop signal sent', 'client_id': client_id, 'stop_key': stop_key }) except Exception as e: logger.error(f"Error stopping VOD client: {e}", exc_info=True) return JsonResponse({'error': str(e)}, status=500) @api_view(["GET"]) @permission_classes([AllowAny]) def stream_xc_movie(request, username, password, stream_id, extension): if not network_access_allowed(request, "STREAMS"): return JsonResponse({"error": "Forbidden"}, status=403) from apps.vod.models import M3UMovieRelation session_id = request.GET.get('session_id') profile_id = request.GET.get('profile_id') user = get_object_or_404(User, username=username) custom_properties = user.custom_properties or {} if "xc_password" not in custom_properties: return Response({"error": "Invalid credentials"}, status=401) if custom_properties["xc_password"] != password: return Response({"error": "Invalid credentials"}, status=401) qs = filter_m3u_relation_queryset_for_playback(M3UMovieRelation.objects.filter(movie_id=stream_id)) qs, _user_pref = filter_relation_qs_for_user_playback(qs, user) try: relations = list(qs.select_related('movie', 'm3u_account')) if not relations: return JsonResponse({"error": "Movie not found"}, status=404) # Pick by least-loaded, then user-preferred, catalog tie-break, then priority from apps.vod.models import M3UMovieRelation as _MMR # noqa movie_relation = _pick_relation_by_account_load(qs, prefer_m3u_account_id=get_vod_catalog_account_id(), user=user) if not movie_relation: return JsonResponse({"error": "Movie not found"}, status=404) except (M3UMovieRelation.DoesNotExist, M3UMovieRelation.MultipleObjectsReturned): return JsonResponse({"error": "Movie not found"}, status=404) return stream_vod(request._request, 'movie', movie_relation.movie.uuid, session_id, profile_id, user) @api_view(["GET"]) @permission_classes([AllowAny]) def stream_xc_episode(request, username, password, stream_id, extension): if not network_access_allowed(request, "STREAMS"): return JsonResponse({"error": "Forbidden"}, status=403) from apps.vod.models import M3UEpisodeRelation session_id = request.GET.get('session_id') profile_id = request.GET.get('profile_id') user = get_object_or_404(User, username=username) custom_properties = user.custom_properties or {} if "xc_password" not in custom_properties: return Response({"error": "Invalid credentials"}, status=401) if custom_properties["xc_password"] != password: return Response({"error": "Invalid credentials"}, status=401) qs = filter_m3u_relation_queryset_for_playback(M3UEpisodeRelation.objects.filter(episode_id=stream_id)) qs, _user_pref = filter_relation_qs_for_user_playback(qs, user) try: episode_relation = _pick_relation_by_account_load(qs, prefer_m3u_account_id=get_vod_catalog_account_id(), user=user) if not episode_relation: return JsonResponse({"error": "Episode not found"}, status=404) except M3UEpisodeRelation.DoesNotExist: return JsonResponse({"error": "Episode not found"}, status=404) return stream_vod(request._request, 'episode', episode_relation.episode.uuid, session_id, profile_id, user) @api_view(["GET"]) @permission_classes([IsAdmin]) def vod_monitor(request): """Combined monitor view: M3U accounts/profiles + Redis profile_connections + active VOD sessions + live channels. Computes REAL connections per profile from vod_persistent_connection:* and ts_proxy:channel:*:clients (via per-client hashes that hold m3u_profile_id). Optional `?reset_zombies=1` resets the Redis profile counters down to the real value when the counter is higher than what we can find in active sessions. Returns JSON consumable by Gestion Dispatcharr UI. """ try: from apps.m3u.models import M3UAccount, M3UAccountProfile from apps.channels.models import Channel from apps.vod.catalog_settings import ( get_vod_catalog_account_id, get_vod_playback_allowed_account_ids, ) connection_manager = MultiWorkerVODConnectionManager.get_instance() redis_client = connection_manager.redis_client catalog_account_id = get_vod_catalog_account_id() playback_allowlist = get_vod_playback_allowed_account_ids() # --- Compute REAL connections per profile_id from active sessions --- real_per_profile: dict[int, int] = {} if redis_client: try: for key in redis_client.scan_iter(match="vod_persistent_connection:*", count=500): data = redis_client.hgetall(key) if not data: continue pid = data.get("m3u_profile_id") try: pid = int(pid) if pid is not None else None except (TypeError, ValueError): pid = None if pid is None: continue real_per_profile[pid] = real_per_profile.get(pid, 0) + 1 except Exception: pass try: for key in redis_client.scan_iter(match="ts_proxy:channel:*:clients:*", count=500): cdata = redis_client.hgetall(key) if not cdata: continue pid = cdata.get("m3u_profile_id") or cdata.get("profile_id") try: pid = int(pid) if pid is not None else None except (TypeError, ValueError): pid = None if pid is None: continue real_per_profile[pid] = real_per_profile.get(pid, 0) + 1 except Exception: pass reset_zombies = request.GET.get("reset_zombies") in ("1", "true", "yes") zombie_actions: list[dict] = [] # Accounts + profiles + connection counts (with reconciliation) accounts_payload = [] for acc in M3UAccount.objects.filter(is_active=True).order_by("id"): cp = acc.custom_properties or {} profiles = [] total_used = 0 total_max = 0 unlimited = False for prof in M3UAccountProfile.objects.filter(m3u_account=acc, is_active=True).order_by("id"): redis_counter = 0 if redis_client: try: redis_counter = int(redis_client.get(f"profile_connections:{prof.id}") or 0) except Exception: redis_counter = 0 real_count = int(real_per_profile.get(int(prof.id), 0)) zombie_delta = max(0, redis_counter - real_count) if zombie_delta > 0 and reset_zombies and redis_client: try: redis_client.set(f"profile_connections:{prof.id}", real_count) zombie_actions.append({ "profile_id": prof.id, "old": redis_counter, "new": real_count, }) redis_counter = real_count zombie_delta = 0 except Exception: pass effective = real_count if zombie_delta > 0 else redis_counter total_used += effective if prof.max_streams == 0: unlimited = True else: total_max += prof.max_streams profiles.append( { "id": prof.id, "name": prof.name, "is_default": prof.is_default, "max_streams": prof.max_streams, "current_connections": effective, "redis_counter": redis_counter, "real_active": real_count, "zombie_delta": zombie_delta, } ) accounts_payload.append( { "id": acc.id, "name": acc.name, "type": acc.account_type, "priority": acc.priority, "enable_vod": bool(cp.get("enable_vod", False)), "is_catalog": (catalog_account_id is not None and acc.id == catalog_account_id), "in_playback_allowlist": (playback_allowlist is None) or (acc.id in playback_allowlist), "profiles": profiles, "total_used": total_used, "total_max": 0 if unlimited else total_max, "unlimited": unlimited, } ) # Active VOD sessions vod_sessions = [] if redis_client: cursor = 0 current_time = time.time() while True: cursor, keys = redis_client.scan(cursor, match="vod_persistent_connection:*", count=200) for key in keys: try: data = redis_client.hgetall(key) if not data: continue key_str = key.decode() if isinstance(key, bytes) else str(key) session_id = key_str.split(":", 1)[1] if ":" in key_str else key_str ctype = data.get("content_obj_type") or "unknown" cuuid = data.get("content_uuid") or "" title = "Unknown" try: if ctype == "movie" and cuuid: m = Movie.objects.filter(uuid=cuuid).only("name").first() if m: title = m.name elif ctype == "episode" and cuuid: e = Episode.objects.select_related("series").filter(uuid=cuuid).first() if e: title = f"{e.series.name} - S{(e.season_number or 0):02d}E{(e.episode_number or 0):02d} - {e.name}" except Exception: pass prof_id = data.get("m3u_profile_id") prof_info = {} try: if prof_id: prof = M3UAccountProfile.objects.select_related("m3u_account").get(id=int(prof_id)) prof_info = { "profile_id": prof.id, "profile_name": prof.name, "account_id": prof.m3u_account_id, "account_name": prof.m3u_account.name, } except Exception: pass try: created_at = float(data.get("created_at") or 0) duration = int(current_time - created_at) if created_at else 0 except Exception: duration = 0 vod_sessions.append( { "session_id": session_id, "content_type": ctype, "content_uuid": cuuid, "title": title, "client_ip": data.get("client_ip") or "", "user_agent": (data.get("client_user_agent") or "")[:160], "user_id": data.get("user_id") or "", "created_at": data.get("created_at") or "", "last_activity": data.get("last_activity") or "", "duration_seconds": duration, "bytes_sent": int(data.get("bytes_sent") or 0), **prof_info, } ) except Exception as exc: logger.debug(f"vod_monitor session parse failed: {exc}") if cursor == 0: break # Live (TS proxy) clients per channel live_clients = [] if redis_client: cursor = 0 channel_clients = {} while True: cursor, keys = redis_client.scan(cursor, match="ts_proxy:channel:*:clients", count=200) for key in keys: key_str = key.decode() if isinstance(key, bytes) else str(key) try: parts = key_str.split(":") channel_id = parts[2] except Exception: continue try: count = int(redis_client.scard(key) or 0) except Exception: count = 0 if count <= 0: continue client_ids = [] try: client_ids = list(redis_client.smembers(key) or [])[:25] client_ids = [c.decode() if isinstance(c, bytes) else c for c in client_ids] except Exception: pass sample_clients = [] for cid in client_ids: try: cdata = redis_client.hgetall(f"ts_proxy:channel:{channel_id}:clients:{cid}") if not cdata: continue sample_clients.append( { "client_id": cid, "ip": cdata.get("ip_address") or cdata.get("ip") or "", "user_agent": (cdata.get("user_agent") or "")[:120], "username": cdata.get("username") or "", "connected_since": cdata.get("connected_since") or cdata.get("connect_time") or "", } ) except Exception: continue channel_clients[channel_id] = {"count": count, "clients": sample_clients} if cursor == 0: break channel_ids = [int(cid) for cid in channel_clients.keys() if str(cid).isdigit()] channel_names = {str(c.id): c.name for c in Channel.objects.filter(id__in=channel_ids).only("id", "name")} for cid, info in channel_clients.items(): live_clients.append( { "channel_id": cid, "channel_name": channel_names.get(str(cid), f"Canal {cid}"), "client_count": info["count"], "clients": info["clients"], } ) return JsonResponse( { "timestamp": time.time(), "catalog_account_id": catalog_account_id, "playback_allowed_account_ids": playback_allowlist, "accounts": accounts_payload, "vod_sessions": vod_sessions, "live_clients": live_clients, "zombies_reset": zombie_actions if reset_zombies else None, } ) except Exception as exc: logger.error(f"vod_monitor error: {exc}", exc_info=True) return JsonResponse({"error": str(exc)}, status=500)