Files
root f56b088643
Base Image Build / prepare (push) Has been cancelled
Build and Push Multi-Arch Docker Image / build-and-push (push) Has been cancelled
Frontend Tests / test (push) Has been cancelled
Base Image Build / docker (amd64, ubuntu-24.04) (push) Has been cancelled
Base Image Build / docker (arm64, ubuntu-24.04-arm) (push) Has been cancelled
Base Image Build / create-manifest (push) Has been cancelled
Proyecto LCX Dispatcharr multicuenta
2026-05-09 21:24:50 +02:00

1572 lines
76 KiB
Python

"""
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: $<name> -> \g<name>, $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 (?<name>...) 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)