Dispatcharr/apps/proxy/live_proxy/views.py
SergeantPanda 820959e390 fix(live): prevent ghost channel sessions from blocking playback indefinitely
This update resolves an issue where ghost channel sessions stuck in the `initializing` state could block playback requests indefinitely. The initialization process now ensures that the `initializing` state is only written to Redis after acquiring the channel ownership lock, preventing race conditions that could leave sessions without an owner. Additionally, failed initializations are cleaned up immediately, allowing subsequent play requests to retry without hitting a stopping gate. The changes enhance the reliability of channel initialization and playback handling. Tests have been updated to validate these improvements.
2026-07-19 16:08:24 +00:00

1210 lines
51 KiB
Python

import json
import time
import random
import re
import pathlib
from django.db import close_old_connections
from django.http import (
StreamingHttpResponse,
JsonResponse,
HttpResponseRedirect,
HttpResponse,
Http404,
)
from django.views.decorators.csrf import csrf_exempt
from django.shortcuts import get_object_or_404
from .server import ProxyServer
from .channel_status import ChannelStatus, build_live_channel_stats_data
from .output.ts.generator import create_stream_generator
from .output.fmp4.generator import create_fmp4_stream_generator
from .utils import get_client_ip
from .redis_keys import RedisKeys
from apps.channels.models import Channel, Stream
from apps.m3u.models import M3UAccount, M3UAccountProfile
from apps.accounts.models import User
from core.models import UserAgent, CoreSettings, PROXY_PROFILE_NAME
from rest_framework.decorators import api_view, permission_classes
from rest_framework.permissions import AllowAny
from rest_framework.response import Response
from apps.accounts.permissions import (
IsAdmin,
permission_classes_by_method,
permission_classes_by_action,
)
from .constants import ChannelState, EventType, StreamType, ChannelMetadataField
from .services.channel_service import ChannelService
from core.utils import send_websocket_update
from .url_utils import (
generate_stream_url,
transform_url,
get_stream_info_for_switch,
get_stream_object,
get_alternate_streams,
)
from .utils import get_logger
from uuid import UUID
import gevent
from dispatcharr.utils import network_access_allowed
from apps.proxy.utils import check_user_stream_limits
logger = get_logger()
def _channel_stopping_response():
response = JsonResponse(
{"error": "Channel is stopping, retry shortly"},
status=503,
)
response["Retry-After"] = "1"
return response
def _channel_setup_needed(proxy_server, channel_id):
"""
Decide whether this worker still needs to run full channel setup.
Returns (needs_setup, state, wait_for_init).
"""
state = None
if proxy_server.redis_client:
metadata = proxy_server.redis_client.hgetall(RedisKeys.channel_metadata(channel_id))
if metadata:
state = metadata.get(ChannelMetadataField.STATE)
if state in (
ChannelState.ACTIVE,
ChannelState.WAITING_FOR_CLIENTS,
ChannelState.BUFFERING,
ChannelState.INITIALIZING,
ChannelState.CONNECTING,
):
wait_for_init = state in (
ChannelState.INITIALIZING,
ChannelState.CONNECTING,
)
return False, state, wait_for_init
if state == ChannelState.STOPPING:
return False, state, False
if state in (ChannelState.ERROR, ChannelState.STOPPED):
return True, state, False
# Unknown/empty state: trust a live owner worker, otherwise re-setup
owner = metadata.get(ChannelMetadataField.OWNER)
if owner:
owner_heartbeat_key = f"live:worker:{owner}:heartbeat"
if proxy_server.redis_client.exists(owner_heartbeat_key):
return False, state, False
return True, state, False
if proxy_server.check_if_channel_exists(channel_id):
return False, state, False
return True, state, False
def _drop_pre_registered_client(proxy_server, channel_id, client_id):
"""Undo an early add_client() when setup aborts before streaming starts."""
mgr = proxy_server.client_managers.get(channel_id)
if mgr:
mgr.remove_client(client_id)
return
if not proxy_server.redis_client:
return
proxy_server.redis_client.srem(RedisKeys.clients(channel_id), client_id)
proxy_server.redis_client.delete(RedisKeys.client_metadata(channel_id, client_id))
def _resolve_output_format(user, force=None, request=None):
"""Return the output format string to use for this client."""
_FORMAT_ALIASES = {
'mpegts': 'mpegts',
'ts': 'mpegts',
'fmp4': 'fmp4',
'mp4': 'fmp4',
}
if force:
return force
if request:
# Support both ?output_format= (native) and ?output= (XC-style)
param = request.GET.get('output_format') or request.GET.get('output')
if param in _FORMAT_ALIASES:
return _FORMAT_ALIASES[param]
if user:
custom = getattr(user, 'custom_properties', None) or {}
user_format = custom.get('output_format')
if user_format:
return user_format
return CoreSettings.get_default_output_format()
def _resolve_output_profile(request, user):
from core.models import OutputProfile
param = request.GET.get('output_profile')
if param:
try:
return OutputProfile.objects.get(id=int(param), is_active=True)
except (OutputProfile.DoesNotExist, ValueError, TypeError):
return None
if user:
custom = getattr(user, 'custom_properties', None) or {}
profile_id = custom.get('output_profile')
if profile_id:
try:
return OutputProfile.objects.get(id=int(profile_id), is_active=True)
except (OutputProfile.DoesNotExist, ValueError, TypeError):
return None
return None
@api_view(["GET"])
@permission_classes([AllowAny])
def stream_ts(request, channel_id, user=None, force_output_format=None):
if not network_access_allowed(request, "STREAMS"):
return JsonResponse({"error": "Forbidden"}, status=403)
"""Stream TS data to client with immediate response and keep-alive packets during initialization"""
if user is None and hasattr(request, 'user') and request.user.is_authenticated:
user = request.user
client_user_agent = None
proxy_server = ProxyServer.get_instance()
connection_allocated = False # Track if connection slot was allocated via get_stream()
# Initialized before the try so the exception handler can always safely
# check/clean it up, regardless of where in the setup a failure occurs.
_client_pre_registered = False
channel = None
client_id = None
channel_display_name = None
try:
channel = get_stream_object(channel_id)
channel_display_name = getattr(channel, "name", None)
# Generate a unique client ID
client_id = f"client_{int(time.time() * 1000)}_{random.randint(1000, 9999)}"
client_ip = get_client_ip(request)
logger.info(f"[{client_id}] Requested stream for channel {channel_id}")
# Extract client user agent early
for header in ["HTTP_USER_AGENT", "User-Agent", "user-agent"]:
if header in request.META:
client_user_agent = request.META[header]
logger.debug(
f"[{client_id}] Client connected with user agent: {client_user_agent}"
)
break
if user:
if not check_user_stream_limits(user, client_id, media_id=channel_id):
return JsonResponse(
{"error": f"Stream limit exceeded ({user.stream_limit} concurrent streams allowed)"},
status=429
)
if ChannelService.is_channel_unavailable_for_new_clients(channel_id):
logger.info(
f"[{client_id}] Channel {channel_id} unavailable. Teardown or pending shutdown"
)
return _channel_stopping_response()
# Check if we need to reinitialize the channel
needs_initialization, channel_state, channel_initializing = _channel_setup_needed(
proxy_server, channel_id
)
if channel_state == ChannelState.STOPPING:
logger.info(
f"[{client_id}] Channel {channel_id} is stopping, rejecting request"
)
return _channel_stopping_response()
if channel_initializing:
logger.debug(
f"[{client_id}] Channel {channel_id} is still initializing, client will wait"
)
elif not needs_initialization:
logger.debug(
f"[{client_id}] Channel {channel_id} in state {channel_state}, skipping initialization"
)
elif channel_state in (ChannelState.ERROR, ChannelState.STOPPED):
logger.info(
f"[{client_id}] Channel {channel_id} in terminal state {channel_state}, will reinitialize"
)
resolved_output_profile = None
resolved_output_format = None
output_options_resolved = False
# Start initialization if needed
if needs_initialization:
if ChannelService.is_channel_unavailable_for_new_clients(channel_id):
logger.info(
f"[{client_id}] Channel {channel_id} became unavailable before init, rejecting"
)
return _channel_stopping_response()
logger.info(f"[{client_id}] Starting channel {channel_id} initialization")
# Force cleanup of any previous instance if in terminal state
if channel_state in [
ChannelState.ERROR,
ChannelState.STOPPING,
ChannelState.STOPPED,
]:
logger.warning(
f"[{client_id}] Channel {channel_id} in state {channel_state}, forcing cleanup"
)
ChannelService.stop_channel(channel_id)
perform_setup = False
owned_for_init = False
init_lock = proxy_server._get_channel_init_lock(channel_id)
init_lock.acquire()
try:
needs_setup, channel_state, wait_for_init = _channel_setup_needed(
proxy_server, channel_id
)
if channel_state == ChannelState.STOPPING or (
ChannelService.is_channel_unavailable_for_new_clients(channel_id)
):
logger.info(
f"[{client_id}] Channel {channel_id} unavailable after init lock, rejecting"
)
return _channel_stopping_response()
if not needs_setup:
if wait_for_init:
channel_initializing = True
logger.info(
f"[{client_id}] Channel {channel_id} already set up after init lock "
f"(state={channel_state}), attaching as follower"
)
elif channel_id in proxy_server._channels_setting_up:
channel_initializing = True
logger.info(
f"[{client_id}] Channel {channel_id} setup already in progress on this "
f"worker, skipping stream reservation and attaching as follower"
)
elif not proxy_server.try_acquire_ownership(channel_id):
channel_initializing = True
logger.info(
f"[{client_id}] Channel {channel_id} owned by another worker, "
f"skipping stream reservation and attaching as follower"
)
else:
owned_for_init = True
proxy_server._channels_setting_up.add(channel_id)
perform_setup = True
finally:
proxy_server._finish_channel_init_lock(channel_id, init_lock)
if perform_setup:
try:
# Use fixed retry interval and timeout
retry_timeout = 3 # 3 seconds total timeout
retry_interval = 0.1 # 100ms between attempts
wait_start_time = time.time()
stream_url = None
stream_user_agent = None
transcode = False
profile_value = None
slot_reserved = False
error_reason = None
attempt = 0
should_retry = True
# Try to get a stream with fixed interval retries
while should_retry and time.time() - wait_start_time < retry_timeout:
attempt += 1
(
stream_url,
stream_user_agent,
transcode,
profile_value,
slot_reserved,
error_reason,
) = generate_stream_url(channel_id)
if stream_url is not None:
logger.info(
f"[{client_id}] Successfully obtained stream for channel {channel_id} after {attempt} attempts"
)
break
# On first failure, check if the error is retryable
if attempt == 1:
if error_reason and "maximum connection limits" not in error_reason:
logger.warning(
f"[{client_id}] Can't retry - error not related to connection limits: {error_reason}"
)
should_retry = False
break
# Check if we have time remaining for another sleep cycle
elapsed_time = time.time() - wait_start_time
remaining_time = retry_timeout - elapsed_time
# If we don't have enough time for the next sleep interval, break
# but only after we've already made an attempt (the while condition will try one more time)
if remaining_time <= retry_interval:
logger.info(
f"[{client_id}] Insufficient time ({remaining_time:.1f}s) for another sleep cycle, will make one final attempt"
)
break
# Wait before retrying
logger.info(
f"[{client_id}] Waiting {retry_interval*1000:.0f}ms for a connection to become available (attempt {attempt}, {remaining_time:.1f}s remaining)"
)
gevent.sleep(retry_interval)
retry_interval += 0.025 # Increase wait time by 25ms for next attempt
# Make one final attempt if we still don't have a stream, should retry, and haven't exceeded timeout
if stream_url is None and should_retry and time.time() - wait_start_time < retry_timeout:
attempt += 1
logger.info(
f"[{client_id}] Making final attempt {attempt} at timeout boundary"
)
(
stream_url,
stream_user_agent,
transcode,
profile_value,
slot_reserved,
error_reason,
) = generate_stream_url(channel_id)
if stream_url is not None:
logger.info(
f"[{client_id}] Successfully obtained stream on final attempt for channel {channel_id}"
)
if stream_url is None:
if slot_reserved and not channel.release_stream():
logger.debug(f"[{client_id}] release_stream found no keys during failed init cleanup")
# Get the specific error message if available
wait_duration = f"{int(time.time() - wait_start_time)}s"
error_msg = (
error_reason
if error_reason
else "No available streams for this channel"
)
logger.info(
f"[{client_id}] Failed to obtain stream after {attempt} attempts over {wait_duration}: {error_msg}"
)
return JsonResponse(
{"error": error_msg, "waited": wait_duration}, status=503
) # 503 Service Unavailable is appropriate here
# generate_stream_url() called get_stream() which allocated a connection
# slot (INCR'd profile_connections) - track this for cleanup on error
if needs_initialization and slot_reserved:
connection_allocated = True
# Read stream assignment from Redis (already set by generate_stream_url → get_stream).
# Avoid calling get_stream() again (INCR profile counter)
# It could double-allocate if the keys were cleared by a concurrent release.
stream_id = None
m3u_profile_id = None
if proxy_server.redis_client:
stream_id_bytes = proxy_server.redis_client.get(f"channel_stream:{channel.id}")
if stream_id_bytes:
stream_id = int(stream_id_bytes)
profile_id_bytes = proxy_server.redis_client.get(f"stream_profile:{stream_id}")
if profile_id_bytes:
m3u_profile_id = int(profile_id_bytes)
logger.info(
f"Channel {channel_id} using stream ID {stream_id}, m3u account profile ID {m3u_profile_id}"
)
# Generate transcode command if needed
stream_profile = channel.get_stream_profile()
if stream_profile.is_redirect():
# Validate the stream URL before redirecting
from .url_utils import (
validate_stream_url,
get_alternate_streams,
get_stream_info_for_switch,
)
# Try initial URL
logger.info(f"[{client_id}] Validating redirect URL: {stream_url}")
is_valid, final_url, status_code, message = validate_stream_url(
stream_url, user_agent=stream_user_agent, timeout=(5, 5)
)
# If first URL doesn't validate, try alternates
if not is_valid:
logger.warning(
f"[{client_id}] Primary stream URL failed validation: {message}"
)
# Track tried streams to avoid loops
tried_streams = {stream_id}
# Get alternate streams
alternates = get_alternate_streams(channel_id, stream_id)
# Try each alternate until one works
for alt in alternates:
if alt["stream_id"] in tried_streams:
continue
tried_streams.add(alt["stream_id"])
# Get stream info
alt_info = get_stream_info_for_switch(
channel_id, alt["stream_id"]
)
if "error" in alt_info:
logger.warning(
f"[{client_id}] Error getting alternate stream info: {alt_info['error']}"
)
continue
# Validate the alternate URL
logger.info(
f"[{client_id}] Trying alternate stream #{alt['stream_id']}: {alt_info['url']}"
)
is_valid, final_url, status_code, message = validate_stream_url(
alt_info["url"],
user_agent=alt_info["user_agent"],
timeout=(5, 5),
)
if is_valid:
logger.info(
f"[{client_id}] Alternate stream #{alt['stream_id']} validated successfully"
)
break
else:
logger.warning(
f"[{client_id}] Alternate stream #{alt['stream_id']} failed validation: {message}"
)
# Release stream lock before redirecting only if we reserved a slot
if connection_allocated and not channel.release_stream():
logger.warning(f"[{client_id}] Failed to release stream before redirect")
connection_allocated = False
# Final decision based on validation results
if is_valid:
logger.info(
f"[{client_id}] Redirecting to validated URL: {final_url} ({message})"
)
# For non-HTTP protocols (RTSP/RTP/UDP), we need to manually create the redirect
# because Django's HttpResponseRedirect blocks them for security
if final_url.startswith(('rtsp://', 'rtp://', 'udp://')):
logger.info(f"[{client_id}] Using manual redirect for non-HTTP protocol")
response = HttpResponse(status=301)
response['Location'] = final_url
return response
return HttpResponseRedirect(final_url)
else:
logger.error(
f"[{client_id}] All available redirect URLs failed validation"
)
return JsonResponse(
{"error": "All available streams failed validation"}, status=502
) # 502 Bad Gateway
# Initialize channel with the stream's user agent (not the client's)
if ChannelService.is_channel_unavailable_for_new_clients(channel_id):
if connection_allocated:
if not channel.release_stream():
logger.warning(f"[{client_id}] Failed to release stream before teardown reject")
connection_allocated = False
logger.info(
f"[{client_id}] Channel {channel_id} unavailable before init call, rejecting"
)
return _channel_stopping_response()
success = ChannelService.initialize_channel(
channel_id,
stream_url,
stream_user_agent,
transcode,
profile_value,
stream_id,
m3u_profile_id,
channel_name=channel.name,
)
if not success:
if connection_allocated:
if not channel.release_stream():
logger.warning(f"[{client_id}] Failed to release stream after init failure")
connection_allocated = False
return JsonResponse(
{"error": "Failed to initialize channel"}, status=500
)
# Channel initialized: lifecycle owns the connection and ownership lock
connection_allocated = False
owned_for_init = False
# If we're the owner, register the client now so the watchdog
# doesn't stop the channel during connection (which can take
# longer than the grace period). The generator handles waiting
# with keepalive packets via _wait_for_initialization().
if proxy_server.am_i_owner(channel_id):
resolved_output_profile = _resolve_output_profile(request, user)
resolved_output_format = _resolve_output_format(user, force_output_format, request)
output_options_resolved = True
resolved_format = (
f'{resolved_output_format}:p{resolved_output_profile.id}'
if resolved_output_profile else resolved_output_format
)
client_manager = proxy_server.client_managers[channel_id]
if not client_manager.add_client(
client_id, client_ip, client_user_agent, user,
output_format=resolved_output_format,
output_profile_id=resolved_output_profile.id if resolved_output_profile else None,
):
logger.error(
f"[{client_id}] Failed to register client with channel {channel_id} during init"
)
return JsonResponse(
{"error": "Failed to register client"}, status=503
)
logger.info(
f"[{client_id}] Client registered with channel {channel_id} "
f"(output: {resolved_format}, profile: {resolved_output_profile.id if resolved_output_profile else None})"
)
_client_pre_registered = True
logger.info(f"[{client_id}] Successfully initialized channel {channel_id}")
channel_initializing = True
finally:
proxy_server._clear_channel_setting_up(channel_id)
if owned_for_init:
proxy_server.release_ownership(channel_id, signal_stopping=False)
# Register client - can do this regardless of initialization state
# Create local resources if needed
if (
channel_id not in proxy_server.stream_buffers
or channel_id not in proxy_server.client_managers
):
logger.debug(
f"[{client_id}] Channel {channel_id} exists in Redis but not initialized in this worker - initializing now"
)
# Get URL from Redis metadata
url = None
stream_user_agent = None # Initialize the variable
if proxy_server.redis_client:
metadata_key = RedisKeys.channel_metadata(channel_id)
url_bytes, ua_bytes, profile_bytes = proxy_server.redis_client.hmget(
metadata_key,
ChannelMetadataField.URL,
ChannelMetadataField.USER_AGENT,
ChannelMetadataField.STREAM_PROFILE,
)
if url_bytes:
url = url_bytes
if ua_bytes:
stream_user_agent = ua_bytes
# Extract transcode setting from Redis
if profile_bytes:
profile_str = profile_bytes
use_transcode = (
profile_str == PROXY_PROFILE_NAME or profile_str == "None"
)
logger.debug(
f"Using profile '{profile_str}' for channel {channel_id}, transcode={use_transcode}"
)
else:
# Default settings when profile not found in Redis
profile_str = "None" # Default profile name
use_transcode = (
False # Default to direct streaming without transcoding
)
logger.debug(
f"No profile found in Redis for channel {channel_id}, defaulting to transcode={use_transcode}"
)
# Use client_user_agent as fallback if stream_user_agent is None
success = proxy_server.initialize_channel(
url,
channel_id,
stream_user_agent or client_user_agent,
use_transcode,
channel_name=channel_display_name,
)
if not success:
logger.error(
f"[{client_id}] Failed to initialize channel {channel_id} locally"
)
return JsonResponse(
{"error": "Failed to initialize channel locally"}, status=500
)
logger.info(
f"[{client_id}] Successfully initialized channel {channel_id} locally"
)
if ChannelService.is_channel_unavailable_for_new_clients(channel_id):
if _client_pre_registered:
_drop_pre_registered_client(proxy_server, channel_id, client_id)
logger.info(
f"[{client_id}] Channel {channel_id} became unavailable during setup, rejecting"
)
return _channel_stopping_response()
if not output_options_resolved:
resolved_output_profile = _resolve_output_profile(request, user)
resolved_output_format = _resolve_output_format(user, force_output_format, request)
# When an output profile is active, append :p{id} to the format key so each
# (format, profile) pair gets its own independent remux pipeline in Redis.
resolved_format = (
f'{resolved_output_format}:p{resolved_output_profile.id}'
if resolved_output_profile else resolved_output_format
)
# Pre-register before slow setup (ensure_output_profile) so the non-owner
# cleanup thread does not tear down local resources while connecting.
if not _client_pre_registered:
client_manager = proxy_server.client_managers.get(channel_id)
if not client_manager:
logger.error(
f"[{client_id}] Channel {channel_id} missing client_manager during setup"
)
return JsonResponse(
{"error": "Channel resources unavailable"}, status=503
)
if not client_manager.add_client(
client_id, client_ip, client_user_agent, user,
output_format=resolved_output_format,
output_profile_id=resolved_output_profile.id if resolved_output_profile else None,
):
logger.error(
f"[{client_id}] Failed to register client with channel {channel_id}"
)
return JsonResponse(
{"error": "Failed to register client"}, status=503
)
_client_pre_registered = True
logger.info(
f"[{client_id}] Client registered with channel {channel_id} "
f"(output: {resolved_format}, profile: {resolved_output_profile.id if resolved_output_profile else None})"
)
if resolved_output_profile:
cmd = resolved_output_profile.build_command()
if not proxy_server.ensure_output_profile(channel_id, resolved_output_profile.id, cmd):
if _client_pre_registered:
_drop_pre_registered_client(proxy_server, channel_id, client_id)
return JsonResponse(
{"error": "Failed to start output profile transcode"}, status=500
)
source_buffer = proxy_server.get_buffer(
channel_id,
profile=resolved_output_profile.id if resolved_output_profile else None
)
client_manager = proxy_server.client_managers.get(channel_id)
if not client_manager:
if _client_pre_registered:
_drop_pre_registered_client(proxy_server, channel_id, client_id)
logger.error(
f"[{client_id}] Channel {channel_id} client_manager removed during setup"
)
return JsonResponse(
{"error": "Channel resources unavailable"}, status=503
)
if resolved_output_format == 'fmp4':
if not proxy_server.ensure_output_format(
channel_id, resolved_format,
source_buffer=source_buffer if resolved_output_profile else None,
):
if _client_pre_registered:
_drop_pre_registered_client(proxy_server, channel_id, client_id)
return JsonResponse(
{"error": "Failed to start output format remux"}, status=500
)
generate = create_fmp4_stream_generator(
channel_id, client_id, client_ip, client_user_agent, channel_initializing, user=user,
fmt=resolved_format,
channel_name=channel_display_name,
)
content_type = "video/mp4"
else:
generate = create_stream_generator(
channel_id,
client_id,
client_ip,
client_user_agent,
channel_initializing,
user=user,
buffer=source_buffer,
channel_name=channel_display_name,
)
content_type = "video/mp2t"
response = StreamingHttpResponse(
streaming_content=generate(), content_type=content_type
)
response["Cache-Control"] = "no-cache"
return response
except Http404:
raise
except Exception as e:
logger.error(f"Error in stream_ts: {e}", exc_info=True)
if connection_allocated and channel is not None:
try:
if not channel.release_stream():
logger.warning(f"[{client_id}] Failed to release stream in exception handler")
except Exception:
pass
# Client may have been pre-registered (before ensure_output_profile /
# get_buffer / generator setup) to protect against the non-owner
# cleanup thread. If setup then failed with an unhandled exception,
# remove it so it doesn't linger as a phantom connection.
if _client_pre_registered:
try:
_drop_pre_registered_client(proxy_server, channel_id, client_id)
except Exception:
logger.warning(f"[{client_id}] Failed to remove client during exception cleanup")
return JsonResponse({"error": str(e)}, status=500)
finally:
# Runs before StreamingHttpResponse is handed to the WSGI server, so the
# request greenlet does not hold a pool slot for the life of the stream.
# Also covers Http404 from get_stream_object (re-raised above).
close_old_connections()
@api_view(["GET"])
@permission_classes([AllowAny])
def stream_xc(request, username, password, channel_id):
try:
user = get_object_or_404(User, username=username)
extension = pathlib.Path(channel_id).suffix
channel_id = pathlib.Path(channel_id).stem
if not network_access_allowed(request, 'STREAMS', user):
return Response({"error": "Forbidden"}, status=403)
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)
if user.user_level < 10:
user_profile_count = user.channel_profiles.count()
# If user has ALL profiles or NO profiles, give unrestricted access
if user_profile_count == 0:
# No profile filtering - user sees all channels based on user_level
filters = {
"id": int(channel_id),
"user_level__lte": user.user_level
}
channel = Channel.objects.filter(**filters).first()
else:
# User has specific limited profiles assigned
filters = {
"id": int(channel_id),
"channelprofilemembership__enabled": True,
"user_level__lte": user.user_level,
"channelprofilemembership__channel_profile__in": user.channel_profiles.all()
}
channel = Channel.objects.filter(**filters).distinct().first()
if not channel:
return JsonResponse({"error": "Not found"}, status=404)
else:
channel = get_object_or_404(Channel, id=channel_id)
if extension.lower() == '.mp4':
force_format = 'fmp4'
elif extension.lower() == '.ts':
force_format = 'mpegts'
else:
force_format = None
return stream_ts(request._request, str(channel.uuid), user, force_output_format=force_format)
except Http404:
raise
finally:
# Auth/channel lookup ORM above; stream_ts also releases on its own paths.
close_old_connections()
@csrf_exempt
@api_view(["POST"])
@permission_classes([IsAdmin])
def change_stream(request, channel_id):
"""Change stream URL for existing channel with enhanced diagnostics"""
proxy_server = ProxyServer.get_instance()
try:
data = json.loads(request.body)
new_url = data.get("url")
user_agent = data.get("user_agent")
stream_id = data.get("stream_id")
m3u_profile_id = None
stream_name = None
# If stream_id is provided, get the URL and user_agent from it
if stream_id:
logger.info(
f"Stream ID {stream_id} provided, looking up stream info for channel {channel_id}"
)
stream_info = get_stream_info_for_switch(channel_id, stream_id)
if "error" in stream_info:
return JsonResponse(
{"error": stream_info["error"], "stream_id": stream_id}, status=404
)
# Use the info from the stream
new_url = stream_info["url"]
user_agent = stream_info["user_agent"]
m3u_profile_id = stream_info.get("m3u_profile_id")
stream_name = stream_info.get("stream_name")
elif not new_url:
return JsonResponse(
{"error": "Either url or stream_id must be provided"}, status=400
)
logger.info(
f"Attempting to change stream for channel {channel_id} to {new_url}"
)
# Use the service layer instead of direct implementation
# Pass stream_id to ensure proper connection tracking
result = ChannelService.change_stream_url(
channel_id, new_url, user_agent, stream_id, m3u_profile_id, stream_name=stream_name
)
# Get the stream manager before updating URL
stream_manager = proxy_server.stream_managers.get(channel_id)
# If we have a stream manager, reset its tried_stream_ids when manually changing streams
if stream_manager:
# Reset tried streams when manually switching URL via API
stream_manager.tried_stream_ids = set()
logger.debug(
f"Reset tried stream IDs for channel {channel_id} during manual stream change"
)
if result.get("status") == "error":
return JsonResponse(
{
"error": result.get("message", "Unknown error"),
"diagnostics": result.get("diagnostics", {}),
},
status=404,
)
if result.get("success") is False:
error_data = {
"error": result.get("message", result.get("error", "Stream switch failed")),
"channel": channel_id,
"url": new_url,
"owner": result.get("direct_update", False),
"worker_id": proxy_server.worker_id,
}
if stream_id:
error_data["stream_id"] = stream_id
# confirmed=False means owner never responded (504); owner reported failure (502)
status_code = 504 if result.get("confirmed") is False else 502
return JsonResponse(error_data, status=status_code)
# Format response based on whether it was a direct update or event-based
response_data = {
"message": "Stream changed successfully",
"channel": channel_id,
"url": new_url,
"owner": result.get("direct_update", False),
"worker_id": proxy_server.worker_id,
}
# Include stream_id in response if it was used
if stream_id:
response_data["stream_id"] = stream_id
return JsonResponse(response_data)
except json.JSONDecodeError:
return JsonResponse({"error": "Invalid JSON"}, status=400)
except Exception as e:
logger.error(f"Failed to change stream: {e}", exc_info=True)
return JsonResponse({"error": str(e)}, status=500)
@api_view(["GET"])
@permission_classes([IsAdmin])
def channel_status(request, channel_id=None):
"""
Returns status information about channels with detail level based on request:
- /status/ returns basic summary of all channels
- /status/{channel_id} returns detailed info about specific channel
"""
proxy_server = ProxyServer.get_instance()
try:
# Check if Redis is available
if not proxy_server.redis_client:
return JsonResponse({"error": "Redis connection not available"}, status=500)
# Handle single channel or all channels
if channel_id:
# Detailed info for specific channel
channel_info = ChannelStatus.get_detailed_channel_info(channel_id)
if channel_info:
return JsonResponse(channel_info)
else:
return JsonResponse(
{"error": f"Channel {channel_id} not found"}, status=404
)
else:
live_stats = build_live_channel_stats_data(proxy_server.redis_client)
# Send WebSocket update with the stats
# Format it the same way the original Celery task did
send_websocket_update(
"updates",
"update",
{
"success": True,
"type": "channel_stats",
"stats": json.dumps(live_stats),
}
)
return JsonResponse(live_stats)
except Exception as e:
logger.error(f"Error in channel_status: {e}", exc_info=True)
return JsonResponse({"error": str(e)}, status=500)
finally:
close_old_connections()
@csrf_exempt
@api_view(["POST", "DELETE"])
@permission_classes([IsAdmin])
def stop_channel(request, channel_id):
"""Stop a channel and release all associated resources using PubSub events"""
try:
logger.info(f"Request to stop channel {channel_id} received")
# Use the service layer instead of direct implementation
result = ChannelService.stop_channel(channel_id)
if result.get("status") == "error":
return JsonResponse(
{"error": result.get("message", "Unknown error")}, status=404
)
return JsonResponse(
{
"message": "Channel stop request sent",
"channel_id": channel_id,
"previous_state": result.get("previous_state"),
}
)
except Exception as e:
logger.error(f"Failed to stop channel: {e}", exc_info=True)
return JsonResponse({"error": str(e)}, status=500)
@csrf_exempt
@api_view(["POST"])
@permission_classes([IsAdmin])
def stop_client(request, channel_id):
"""Stop a specific client connection using existing client management"""
try:
# Parse request body to get client ID
data = json.loads(request.body)
client_id = data.get("client_id")
if not client_id:
return JsonResponse({"error": "No client_id provided"}, status=400)
# Use the service layer instead of direct implementation
result = ChannelService.stop_client(channel_id, client_id)
if result.get("status") == "error":
return JsonResponse({"error": result.get("message")}, status=404)
return JsonResponse(
{
"message": "Client stop request processed",
"channel_id": channel_id,
"client_id": client_id,
"locally_processed": result.get("locally_processed", False),
}
)
except json.JSONDecodeError:
return JsonResponse({"error": "Invalid JSON"}, status=400)
except Exception as e:
logger.error(f"Failed to stop client: {e}", exc_info=True)
return JsonResponse({"error": str(e)}, status=500)
@csrf_exempt
@api_view(["POST"])
@permission_classes([IsAdmin])
def next_stream(request, channel_id):
"""Switch to the next available stream for a channel"""
proxy_server = ProxyServer.get_instance()
try:
logger.info(
f"Request to switch to next stream for channel {channel_id} received"
)
# Check if the channel exists
channel = get_stream_object(channel_id)
# First check if channel is active in Redis
current_stream_id = None
profile_id = None
if proxy_server.redis_client:
metadata_key = RedisKeys.channel_metadata(channel_id)
if proxy_server.redis_client.exists(metadata_key):
# Get current stream ID from Redis
stream_id_bytes = proxy_server.redis_client.hget(
metadata_key, ChannelMetadataField.STREAM_ID
)
if stream_id_bytes:
current_stream_id = int(stream_id_bytes)
logger.info(
f"Found current stream ID {current_stream_id} in Redis for channel {channel_id}"
)
# Get M3U profile from Redis if available
profile_id_bytes = proxy_server.redis_client.hget(
metadata_key, ChannelMetadataField.M3U_PROFILE
)
if profile_id_bytes:
profile_id = int(profile_id_bytes)
logger.info(
f"Found M3U profile ID {profile_id} in Redis for channel {channel_id}"
)
if not current_stream_id:
# Channel is not running
return JsonResponse(
{"error": "No current stream found for channel"}, status=404
)
# Get all streams for this channel in their defined order
streams = list(channel.streams.all().order_by("channelstream__order"))
if len(streams) <= 1:
return JsonResponse(
{
"error": "No alternate streams available for this channel",
"current_stream_id": current_stream_id,
},
status=404,
)
# Find the current stream's position in the list
current_index = None
for i, stream in enumerate(streams):
if stream.id == current_stream_id:
current_index = i
break
if current_index is None:
logger.warning(
f"Current stream ID {current_stream_id} not found in channel's streams list"
)
# Fall back to the first stream that's not the current one
next_stream = next((s for s in streams if s.id != current_stream_id), None)
if not next_stream:
return JsonResponse(
{
"error": "Could not find current stream in channel list",
"current_stream_id": current_stream_id,
},
status=404,
)
else:
# Get the next stream in the rotation (with wrap-around)
next_index = (current_index + 1) % len(streams)
next_stream = streams[next_index]
next_stream_id = next_stream.id
logger.info(
f"Rotating to next stream ID {next_stream_id} for channel {channel_id}"
)
# Get full stream info including URL for the next stream
stream_info = get_stream_info_for_switch(channel_id, next_stream_id)
if "error" in stream_info:
return JsonResponse(
{
"error": stream_info["error"],
"current_stream_id": current_stream_id,
"next_stream_id": next_stream_id,
},
status=404,
)
# Now use the ChannelService to change the stream URL
result = ChannelService.change_stream_url(
channel_id,
stream_info["url"],
stream_info["user_agent"],
next_stream_id,
stream_info.get("m3u_profile_id"),
stream_name=stream_info.get("stream_name"),
)
if result.get("status") == "error":
return JsonResponse(
{
"error": result.get("message", "Unknown error"),
"diagnostics": result.get("diagnostics", {}),
"current_stream_id": current_stream_id,
"next_stream_id": next_stream_id,
},
status=404,
)
if result.get("success") is False:
return JsonResponse(
{
"error": result.get("message", result.get("error", "Stream switch failed")),
"current_stream_id": current_stream_id,
"next_stream_id": next_stream_id,
"owner": result.get("direct_update", False),
"worker_id": proxy_server.worker_id,
},
status=504 if result.get("confirmed") is False else 502,
)
# Format success response
response_data = {
"message": "Stream switched to next available",
"channel": channel_id,
"previous_stream_id": current_stream_id,
"new_stream_id": next_stream_id,
"new_url": stream_info["url"],
"owner": result.get("direct_update", False),
"worker_id": proxy_server.worker_id,
}
return JsonResponse(response_data)
except Http404:
raise
except Exception as e:
logger.error(f"Failed to switch to next stream: {e}", exc_info=True)
return JsonResponse({"error": str(e)}, status=500)
finally:
close_old_connections()