You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
 
 
 
 

192 lines
5.9 KiB

from __future__ import annotations
import logging
from typing import Any, Dict, Optional
import requests
from django.conf import settings
from django.core.exceptions import ImproperlyConfigured
from django.utils import timezone
logger = logging.getLogger(__name__)
class WebEgressError(Exception):
def __init__(
self,
message: str,
*,
status_code: Optional[int] = None,
response_data: Optional[Dict[str, Any]] = None,
):
super().__init__(message)
self.status_code = status_code
self.response_data = response_data or {}
class WebEgressClient:
def __init__(
self,
*,
base_url: Optional[str] = None,
service_token: Optional[str] = None,
timeout: Optional[float] = None,
):
self.base_url = (
base_url
or getattr(settings, "ONLINE_CLASS_WEB_EGRESS_SERVICE_URL", "")
).rstrip("/")
self.service_token = (
service_token
or getattr(settings, "ONLINE_CLASS_WEB_EGRESS_SERVICE_TOKEN", "")
)
self.timeout = timeout or getattr(
settings, "ONLINE_CLASS_WEB_EGRESS_TIMEOUT", 20.0
)
if not self.base_url:
raise ImproperlyConfigured(
"ONLINE_CLASS_WEB_EGRESS_SERVICE_URL must be configured."
)
def start_recording(self, payload: Dict[str, Any]) -> Dict[str, Any]:
return self._request("POST", "/api/web-egress/start/", json=payload)
def stop_recording(self, payload: Dict[str, Any]) -> Dict[str, Any]:
return self._request("POST", "/api/web-egress/stop/", json=payload)
def get_status(self, *, room_id: str, session_id: int) -> Dict[str, Any]:
return self._request(
"GET",
"/api/web-egress/status/",
params={"room_id": room_id, "session_id": session_id},
)
def _request(self, method: str, path: str, **kwargs) -> Dict[str, Any]:
headers = kwargs.pop("headers", {})
headers["Content-Type"] = "application/json"
if self.service_token:
headers["Authorization"] = f"Bearer {self.service_token}"
request_url = f"{self.base_url}{path}"
try:
response = requests.request(
method,
request_url,
headers=headers,
timeout=self.timeout,
**kwargs,
)
except requests.RequestException as exc:
logger.exception(
"[WebEgressClient] Request failed method=%s url=%s error=%s",
method,
request_url,
str(exc),
)
raise WebEgressError("Failed to reach WebEgress service.") from exc
data = self._safe_json(response)
if response.status_code >= 400:
message = (
(data or {}).get("message")
or (data or {}).get("detail")
or response.text
or "WebEgress service request failed."
)
logger.error(
"[WebEgressClient] Non-success response method=%s url=%s status=%s message=%s payload=%s",
method,
request_url,
response.status_code,
message,
data,
)
raise WebEgressError(
message,
status_code=response.status_code,
response_data=data,
)
if data is None:
raise WebEgressError("WebEgress service returned an invalid response.")
if isinstance(data, dict) and data.get("status") is False:
raise WebEgressError(
data.get("message") or data.get("detail") or "WebEgress failed.",
status_code=response.status_code,
response_data=data,
)
return data
@staticmethod
def _safe_json(response: requests.Response) -> Optional[Dict[str, Any]]:
try:
return response.json()
except ValueError:
return None
def is_web_egress_recording_status(status_value: str) -> bool:
return str(status_value or "").strip().lower() in {
"starting",
"recording",
"stopping",
}
def stop_session_web_egress_if_active(session) -> bool:
"""
Best-effort stop for an active WebEgress recording.
Returns True when a stop request was sent, otherwise False.
This helper never raises and is safe to call from close-session flows.
"""
if not getattr(settings, "ONLINE_CLASS_WEB_EGRESS_ENABLED", False):
return False
if not session or not is_web_egress_recording_status(
getattr(session, "web_egress_status", "")
):
return False
try:
client = WebEgressClient()
client.stop_recording(
{
"session_id": session.id,
"room_id": session.room_id,
"egress_id": session.web_egress_id,
}
)
except (ImproperlyConfigured, WebEgressError) as exc:
logger.warning(
"[WebEgress] Failed to stop active recording during session close - session_id=%s room_id=%s error=%s",
session.id,
session.room_id,
str(exc),
)
return False
session.web_egress_status = "stopping"
if not session.web_egress_stopped_at:
session.web_egress_stopped_at = timezone.now()
session.save(
update_fields=[
"web_egress_status",
"web_egress_stopped_at",
"updated_at",
]
)
else:
session.save(update_fields=["web_egress_status", "updated_at"])
logger.info(
"[WebEgress] Stop requested during session close - session_id=%s room_id=%s egress_id=%s",
session.id,
session.room_id,
session.web_egress_id,
)
return True