import json import hmac import hashlib import logging import os import tempfile import subprocess import base64 import jwt from typing import Dict, Any from django.conf import settings from django.core.exceptions import ImproperlyConfigured from django.utils import timezone from django.utils.decorators import method_decorator from django.views.decorators.csrf import csrf_exempt from django.core.files.base import ContentFile from rest_framework import status from rest_framework.views import APIView from rest_framework.response import Response from rest_framework.permissions import AllowAny from rest_framework.parsers import BaseParser from drf_yasg.utils import swagger_auto_schema from drf_yasg import openapi from apps.course.models import CourseLiveSession, LiveSessionUser, Course, Participant, LiveSessionRecording from apps.account.models import User from apps.course.services.plugnmeet import PlugNMeetClient, PlugNMeetError from utils.exceptions import AppAPIException logger = logging.getLogger(__name__) class RawJSONParser(BaseParser): """ Parser that preserves the raw body bytes for HMAC signature verification. """ media_type = '*/*' def parse(self, stream, media_type=None, parser_context=None): return stream.read() @method_decorator(csrf_exempt, name='dispatch') class PlugNMeetWebhookAPIView(APIView): """ Webhook endpoint to receive and process events from the PlugNMeet server. Handles: - room_finished: Closes the live session record. - participant_joined: Tracks student entry. - participant_left: Tracks student exit. - end_recording: Downloads and saves session recordings. """ authentication_classes = [] permission_classes = [AllowAny] parser_classes = [RawJSONParser] @swagger_auto_schema( operation_description="Handle webhook events from PlugNMeet server for live sessions", tags=["Imam-Javad - Course"], responses={ 200: openapi.Response(description="Webhook processed successfully"), 403: openapi.Response(description="Invalid signature"), 400: openapi.Response(description="Invalid payload") } ) def post(self, request, *args, **kwargs): logger.info("⚡ [PlugNMeet Webhook] Request received") # 1. Extract Signature hash_token = request.headers.get('Hash-Token') or request.META.get('HTTP_HASH_TOKEN') if not hash_token: logger.error("❌ [PlugNMeet Webhook] Missing Hash-Token header") return Response({'message': 'Missing Hash-Token header'}, status=403) # 2. Verify Signature if not self._verify_webhook_signature(request, hash_token): return Response({'message': 'Invalid webhook signature'}, status=403) # 3. Parse Payload try: body_bytes = request.data # RawJSONParser puts bytes here payload = json.loads(body_bytes.decode('utf-8')) except Exception as e: logger.error(f"❌ [PlugNMeet Webhook] Parsing Error: {e}") return Response({'message': 'Invalid JSON'}, status=400) event = payload.get('event') logger.info(f"✅ [PlugNMeet Webhook] Event: {event}") # 4. Route Event handler_map = { 'room_finished': self._handle_room_finished, 'participant_joined': self._handle_participant_joined, 'participant_left': self._handle_participant_left, 'recording_proceeded': self._handle_recording_proceeded, } handler = handler_map.get(event) if not handler: logger.info(f"ℹ️ [PlugNMeet Webhook] Event {event} ignored") return Response({'status': 'ok', 'message': f'Event {event} ignored'}) try: result = handler(payload) return Response({'status': 'ok', **result}) except Exception as e: logger.error(f"❌ [PlugNMeet Webhook] Error in {event}: {e}", exc_info=True) return Response({'status': 'error', 'message': str(e)}, status=500) def _verify_webhook_signature(self, request, hash_token: str) -> bool: api_secret = getattr(settings, 'PLUGNMEET_API_SECRET', None) if not api_secret: logger.error("❌ [PlugNMeet Webhook] PLUGNMEET_API_SECRET not configured") return False body_bytes = request.data auth_header = request.headers.get('Authorization') or request.META.get('HTTP_AUTHORIZATION', '') logger.info(f"🔍 [PlugNMeet Webhook] Verify Signature details:") logger.info(f" - Hash-Token (first 30 chars): {hash_token[:30]}... (length: {len(hash_token)})") logger.info(f" - Authorization (first 30 chars): {auth_header[:30]}... (length: {len(auth_header)})") logger.info(f" - Body bytes length: {len(body_bytes)}") if len(body_bytes) > 0: logger.info(f" - Body preview (first 100 bytes): {body_bytes[:100]}") # 1. Try JWT verification (PlugNMeet production standard) try: # A JWT typically consists of three segments separated by dots if len(hash_token.split('.')) == 3: logger.info(" - Attempting JWT decode on Hash-Token...") decoded = jwt.decode( hash_token, api_secret, algorithms=['HS256'], options={ "verify_signature": True, "verify_exp": False, "verify_nbf": False, "verify_iat": False, "verify_aud": False, "verify_iss": False } ) logger.info(f" - Decoded JWT claims: {list(decoded.keys())}") token_sha256 = decoded.get('sha256') if token_sha256: hasher = hashlib.sha256() hasher.update(body_bytes) computed_sha256 = base64.b64encode(hasher.digest()).decode('utf-8') logger.info(f" - Claim SHA256: {token_sha256}") logger.info(f" - Computed SHA256: {computed_sha256}") if hmac.compare_digest(token_sha256, computed_sha256): logger.info(" - JWT SHA256 signature verification successful!") return True else: logger.error(f"❌ [PlugNMeet Webhook] SHA256 mismatch! Token claim: {token_sha256}, Computed: {computed_sha256}") else: logger.error("❌ [PlugNMeet Webhook] JWT missing 'sha256' claim") else: logger.info(" - Hash-Token does not have 3 segments (not a JWT).") except jwt.PyJWTError as e: logger.error(f"❌ [PlugNMeet Webhook] JWT decoding failed: {e.__class__.__name__}: {str(e)}") logger.warning(f"⚠️ [PlugNMeet Webhook] JWT decoding failed. Falling back to HMAC-SHA256 check.") # 2. Fallback: HMAC-SHA256 verification (for backward compatibility with test_webhook.py) expected_signature = hmac.new( api_secret.encode('utf-8'), body_bytes, hashlib.sha256 ).hexdigest() logger.info(f" - Expected HMAC: {expected_signature}") if hmac.compare_digest(hash_token, expected_signature): logger.info(" - Fallback HMAC verification successful!") return True logger.error(f"❌ [PlugNMeet Webhook] Signature mismatch! \nReceived: {hash_token[:10]}...\nExpected HMAC: {expected_signature[:10]}...") return False def _handle_room_finished(self, payload: Dict[str, Any]) -> Dict[str, Any]: room_data = payload.get('room', {}) room_id = room_data.get('room_id') or room_data.get('roomId') or room_data.get('identity') if not room_id: logger.warning(f"⚠️ [PlugNMeet Webhook] Missing room_id in room_finished event. Payload: {payload}") return {'message': 'Missing room identity'} try: session = CourseLiveSession.objects.get(room_id=room_id, ended_at__isnull=True) now = timezone.now() session.ended_at = now session.save(update_fields=['ended_at', 'updated_at']) # Close active user sessions updated_count = LiveSessionUser.objects.filter( session=session, is_online=True, exited_at__isnull=True ).update(is_online=False, exited_at=now, updated_at=now) logger.info(f"🏁 [PlugNMeet Webhook] Session {session.id} ended. Users disconnected: {updated_count}") return {'session_id': session.id, 'closed_users': updated_count} except CourseLiveSession.DoesNotExist: logger.warning(f"⚠️ [PlugNMeet Webhook] room_finished: No active session found for room_id {room_id}") return {'message': 'No active session found for this room'} except Exception as e: logger.error(f"❌ [PlugNMeet Webhook] Error in room_finished for room_id {room_id}: {e}", exc_info=True) return {'error': str(e)} def _handle_participant_joined(self, payload: Dict[str, Any]) -> Dict[str, Any]: room_data = payload.get('room', {}) participant_data = payload.get('participant', {}) room_id = room_data.get('room_id') or room_data.get('roomId') or room_data.get('identity') user_id = participant_data.get('identity') or participant_data.get('user_id') or participant_data.get('userId') if not room_id or not user_id: logger.warning(f"⚠️ [PlugNMeet Webhook] Missing room_id or user_id in participant_joined. room_id: {room_id}, user_id: {user_id}. Payload: {payload}") return {'message': 'Missing required metadata'} try: if not str(user_id).isdigit(): logger.info(f"ℹ️ [PlugNMeet Webhook] Ignoring non-integer user join: {user_id}") return {'message': f'Ignored non-integer user_id {user_id}'} session = CourseLiveSession.objects.get(room_id=room_id, ended_at__isnull=True) user = User.objects.get(id=int(user_id)) role = 'moderator' if user.can_manage_course(session.course) else 'participant' session_user, created = LiveSessionUser.objects.update_or_create( session=session, user=user, defaults={ 'role': role, 'is_online': True, 'exited_at': None, 'entered_at': timezone.now() } ) logger.info(f"👤 [PlugNMeet Webhook] User {user.id} joined session {session.id} (created: {created})") return {'session_user_id': session_user.id, 'created': created} except Exception as e: logger.error(f"❌ [PlugNMeet Webhook] Error in participant_joined (room_id: {room_id}, user_id: {user_id}): {e}", exc_info=True) return {'error': str(e)} def _handle_participant_left(self, payload: Dict[str, Any]) -> Dict[str, Any]: room_data = payload.get('room', {}) participant_data = payload.get('participant', {}) room_id = room_data.get('room_id') or room_data.get('roomId') or room_data.get('identity') user_id = participant_data.get('identity') or participant_data.get('user_id') or participant_data.get('userId') if not room_id or not user_id: logger.warning(f"⚠️ [PlugNMeet Webhook] Missing room_id or user_id in participant_left. room_id: {room_id}, user_id: {user_id}. Payload: {payload}") return {'message': 'Missing required metadata'} try: if not str(user_id).isdigit(): logger.info(f"ℹ️ [PlugNMeet Webhook] Ignoring non-integer user leave: {user_id}") return {'message': f'Ignored non-integer user_id {user_id}'} session = CourseLiveSession.objects.get(room_id=room_id) user = User.objects.get(id=int(user_id)) updated = LiveSessionUser.objects.filter( session=session, user=user, is_online=True ).update(is_online=False, exited_at=timezone.now(), updated_at=timezone.now()) logger.info(f"🚪 [PlugNMeet Webhook] User {user.id} left session {session.id} (updated entries count: {updated})") return {'updated': bool(updated)} except Exception as e: logger.error(f"❌ [PlugNMeet Webhook] Error in participant_left (room_id: {room_id}, user_id: {user_id}): {e}", exc_info=True) return {'error': str(e)} def _get_video_duration(self, video_path: str) -> timezone.timedelta: """Get video duration using ffprobe, return as timedelta.""" try: cmd = [ 'ffprobe', '-v', 'error', '-show_entries', 'format=duration', '-of', 'default=noprint_wrappers=1:nokey=1', video_path ] result = subprocess.run( cmd, stdout=subprocess.PIPE, stderr=subprocess.PIPE, timeout=30 ) if result.returncode == 0: duration_seconds = float(result.stdout.decode().strip()) return timezone.timedelta(seconds=duration_seconds) except Exception as e: logger.warning(f"⚠️ [PlugNMeet Webhook] Duration extraction failed: {e}") return timezone.timedelta(seconds=0) def _handle_recording_proceeded(self, payload: Dict[str, Any]) -> Dict[str, Any]: logger.info(f"📥 [PlugNMeet Webhook] Processing recording_proceeded payload: {payload}") room_data = payload.get('room', {}) recording_info = payload.get('recording_info', {}) room_id = room_data.get('room_id') or room_data.get('roomId') or room_data.get('identity') recording_id = ( recording_info.get('record_id') or recording_info.get('recordId') or recording_info.get('recording_id') or recording_info.get('recordingId') ) file_path = recording_info.get('file_path') or recording_info.get('filePath') if not room_id or not recording_id: print("ERROR------------------------", recording_id, room_id, payload) logger.warning(f"⚠️ [PlugNMeet Webhook] Missing room_id or recording_id in recording_proceeded event. room_id: {room_id}, recording_id: {recording_id}") return {'message': 'Missing recording metadata'} file_name = os.path.basename(file_path) if file_path else f"{recording_id}.mp4" try: session = CourseLiveSession.objects.get(room_id=room_id) client = PlugNMeetClient() # 1. Fetch download token token_response = client.get_recording_download_token(recording_id) if not token_response.get('status'): return {'error': 'Failed to get download token'} download_token = token_response.get('token') download_path = f"/download/recording/{download_token}" # 2. Download to temporary file with tempfile.NamedTemporaryFile(delete=False, suffix='.mp4') as tmp_file: tmp_file_path = tmp_file.name try: logger.info(f"📥 [PlugNMeet Webhook] Downloading recording {recording_id}...") client.download_file(download_path, tmp_file_path) # 3. Read duration using ffprobe duration = self._get_video_duration(tmp_file_path) # 4. Save to Database with open(tmp_file_path, 'rb') as f: content = f.read() recording = LiveSessionRecording.objects.create( session=session, title=f"{session.subject} - Recording", file_time=duration if duration.total_seconds() > 0 else None, recording_type='video' if file_name.lower().endswith('.mp4') else 'voice' ) recording.file.save(file_name, ContentFile(content), save=True) # 5. Generate thumbnail (Optional) self._generate_video_thumbnail(tmp_file_path, recording) logger.info(f"💾 [PlugNMeet Webhook] Recording saved successfully: {recording.id}") return {'recording_id': recording.id, 'file': file_name} finally: if os.path.exists(tmp_file_path): os.unlink(tmp_file_path) except Exception as e: logger.error(f"❌ [PlugNMeet Webhook] RECORDING_PROCEEDED Error: {e}", exc_info=True) return {'error': str(e)} def _generate_video_thumbnail(self, video_path: str, recording: LiveSessionRecording) -> bool: try: with tempfile.NamedTemporaryFile(delete=False, suffix='.jpg') as tmp_thumb: thumbnail_path = tmp_thumb.name cmd = [ 'ffmpeg', '-ss', '1', '-i', video_path, '-frames:v', '1', '-q:v', '2', '-vf', 'scale=640:-1', '-y', thumbnail_path ] result = subprocess.run(cmd, stdout=subprocess.PIPE, stderr=subprocess.PIPE, timeout=30) if result.returncode == 0 and os.path.exists(thumbnail_path) and os.path.getsize(thumbnail_path) > 0: with open(thumbnail_path, 'rb') as f: recording.thumbnail.save(f"thumb_{recording.id}.jpg", ContentFile(f.read()), save=True) os.unlink(thumbnail_path) return True return False except Exception as e: logger.warning(f"⚠️ [PlugNMeet Webhook] Thumbnail failed: {e}") return False