Files
training-software/training/consumers.py
Paperclip CTO dafbfa2bc4 feat(TRA-372): live groupcall backend — WebRTC signaling + session control
- Add Django Channels 4 + channels-redis to requirements and INSTALLED_APPS
- Upgrade ASGI config to ProtocolTypeRouter with JWT-authenticated WebSocket routing
- Add JWTAuthMiddleware for WebSocket token auth via query param
- Add CallSession, CallParticipant, CallEvent models with migration 0004
- MeetingCallConsumer: join/leave, P2P SDP/ICE relay via per-user groups,
  instructor mute/unmute/kick with DB audit and WS broadcast
- REST endpoints: GET/POST/DELETE /meetings/{id}/call/ (session lifecycle),
  POST /moderate/ (mute/unmute/kick), POST /screen-share/, GET /events/ (audit log)
- IsMeetingModerator permission (accepts training:signoff or meeting:moderate)
- Services: get_or_create_call_session, end_call_session, apply_moderation_action,
  toggle_screen_share with full CallEvent audit trail
- Integration tests covering REST lifecycle, moderation, and WebSocket signaling

Co-Authored-By: Paperclip <noreply@paperclip.ing>
2026-05-18 14:12:48 +02:00

362 lines
12 KiB
Python

import json
import logging
from channels.db import database_sync_to_async
from channels.generic.websocket import AsyncWebsocketConsumer
from django.utils.timezone import now
from .models import (
CallEvent,
CallEventType,
CallParticipant,
CallParticipantStatus,
CallSession,
CallSessionStatus,
Meeting,
)
logger = logging.getLogger(__name__)
SIGNALING_MSG_TYPES = {"offer", "answer", "ice_candidate"}
MODERATION_MSG_TYPES = {"mute", "unmute", "kick"}
class MeetingCallConsumer(AsyncWebsocketConsumer):
"""
WebSocket consumer for WebRTC signaling in a meeting call.
URL: /ws/meetings/<meeting_id>/call/
Client sends JSON:
Signaling (P2P relay):
{"type": "offer", "to": "<user_id>", "sdp": {...}}
{"type": "answer", "to": "<user_id>", "sdp": {...}}
{"type": "ice_candidate", "to": "<user_id>", "candidate": {...}}
Moderation (instructor only):
{"type": "mute", "user_id": "<target>"}
{"type": "unmute", "user_id": "<target>"}
{"type": "kick", "user_id": "<target>"}
Server → client events:
{"type": "session_state", "session_id": ..., "participants": [...]}
{"type": "participant_joined","user_id":..., "user_name":..., "is_muted":...}
{"type": "participant_left", "user_id":...}
{"type": "offer"|"answer"|"ice_candidate", "from":..., "sdp"|"candidate":...}
{"type": "muted"|"unmuted", "user_id":..., "by":...}
{"type": "kicked", "user_id":..., "kicked_by":...}
{"type": "session_ended"}
"""
async def connect(self):
user = self.scope.get("user")
if not user or not user.is_authenticated:
await self.close(code=4001)
return
self.meeting_id = self.scope["url_route"]["kwargs"]["meeting_id"]
self.group_name = f"meeting_call_{self.meeting_id}"
# per-user group for targeted P2P signal relay
self.user_group = f"meeting_call_{self.meeting_id}_user_{user.id}"
self.user = user
session, allowed = await self._get_or_validate_session()
if not allowed:
await self.close(code=4003)
return
self.session_id = str(session.id)
await self.channel_layer.group_add(self.group_name, self.channel_name)
await self.channel_layer.group_add(self.user_group, self.channel_name)
await self.accept()
participant = await self._record_join()
await self.channel_layer.group_send(
self.group_name,
{
"type": "call.participant_joined",
"user_id": str(user.id),
"user_name": user.get_full_name() or str(user),
"is_muted": participant.is_muted,
},
)
existing = await self._get_active_participants()
await self.send(text_data=json.dumps({
"type": "session_state",
"session_id": self.session_id,
"participants": existing,
}))
async def disconnect(self, close_code):
if not hasattr(self, "group_name"):
return
await self._record_leave()
await self.channel_layer.group_send(
self.group_name,
{"type": "call.participant_left", "user_id": str(self.user.id)},
)
await self.channel_layer.group_discard(self.group_name, self.channel_name)
await self.channel_layer.group_discard(self.user_group, self.channel_name)
async def receive(self, text_data):
try:
data = json.loads(text_data)
except json.JSONDecodeError:
await self._send_error("invalid_json")
return
msg_type = data.get("type")
if msg_type in SIGNALING_MSG_TYPES:
await self._handle_signaling(msg_type, data)
elif msg_type in MODERATION_MSG_TYPES:
await self._handle_moderation(msg_type, data)
else:
await self._send_error("unknown_message_type")
# --- signaling relay ---
async def _handle_signaling(self, msg_type, data):
to_user_id = str(data.get("to", ""))
if not to_user_id:
await self._send_error("missing_to")
return
target_group = f"meeting_call_{self.meeting_id}_user_{to_user_id}"
payload: dict = {
"type": "call.signal",
"signal_type": msg_type,
"from_user_id": str(self.user.id),
}
if msg_type in ("offer", "answer"):
payload["sdp"] = data.get("sdp")
else:
payload["candidate"] = data.get("candidate")
await self.channel_layer.group_send(target_group, payload)
# --- moderation (instructor only) ---
async def _handle_moderation(self, msg_type, data):
if not await self._is_instructor():
await self._send_error("not_authorized")
return
target_user_id = str(data.get("user_id", ""))
if not target_user_id:
await self._send_error("missing_user_id")
return
if msg_type == "kick":
await self._handle_kick(target_user_id)
elif msg_type in ("mute", "unmute"):
await self._handle_mute_toggle(msg_type, target_user_id)
async def _handle_kick(self, target_user_id):
await self._log_moderation_event(CallEventType.KICK, target_user_id)
await self.channel_layer.group_send(
self.group_name,
{
"type": "call.kicked",
"user_id": target_user_id,
"kicked_by": str(self.user.id),
},
)
async def _handle_mute_toggle(self, action, target_user_id):
is_muted = action == "mute"
await self._set_participant_muted(target_user_id, is_muted)
event_type = CallEventType.MUTE if is_muted else CallEventType.UNMUTE
await self._log_moderation_event(event_type, target_user_id)
broadcast_type = "call.muted" if is_muted else "call.unmuted"
await self.channel_layer.group_send(
self.group_name,
{
"type": broadcast_type,
"user_id": target_user_id,
"by": str(self.user.id),
},
)
# --- channel message handlers (called by channel layer) ---
async def call_participant_joined(self, event):
await self.send(text_data=json.dumps({
"type": "participant_joined",
"user_id": event["user_id"],
"user_name": event["user_name"],
"is_muted": event["is_muted"],
}))
async def call_participant_left(self, event):
await self.send(text_data=json.dumps({
"type": "participant_left",
"user_id": event["user_id"],
}))
async def call_signal(self, event):
payload: dict = {
"type": event["signal_type"],
"from": event["from_user_id"],
}
if event["signal_type"] in ("offer", "answer"):
payload["sdp"] = event.get("sdp")
else:
payload["candidate"] = event.get("candidate")
await self.send(text_data=json.dumps(payload))
async def call_muted(self, event):
await self.send(text_data=json.dumps({
"type": "muted",
"user_id": event["user_id"],
"by": event["by"],
}))
async def call_unmuted(self, event):
await self.send(text_data=json.dumps({
"type": "unmuted",
"user_id": event["user_id"],
"by": event["by"],
}))
async def call_kicked(self, event):
await self.send(text_data=json.dumps({
"type": "kicked",
"user_id": event["user_id"],
"kicked_by": event["kicked_by"],
}))
if event["user_id"] == str(self.user.id):
await self.close(code=4004)
async def call_session_ended(self, event):
await self.send(text_data=json.dumps({"type": "session_ended"}))
await self.close(code=1000)
# --- database helpers ---
@database_sync_to_async
def _get_or_validate_session(self):
try:
meeting = Meeting.objects.get(pk=self.meeting_id)
except Meeting.DoesNotExist:
return None, False
is_participant = meeting.participants.filter(user=self.user).exists()
if not (is_participant or _user_is_instructor(self.user)):
return None, False
session, _ = CallSession.objects.get_or_create(
meeting=meeting,
defaults={"started_by": self.user, "status": CallSessionStatus.ACTIVE},
)
if session.status == CallSessionStatus.ENDED:
return session, False
return session, True
@database_sync_to_async
def _record_join(self):
session = CallSession.objects.get(pk=self.session_id)
participant, created = CallParticipant.objects.get_or_create(
session=session,
user=self.user,
defaults={"status": CallParticipantStatus.JOINED},
)
if not created:
participant.status = CallParticipantStatus.JOINED
participant.is_muted = False
participant.left_at = None
participant.save(update_fields=["status", "is_muted", "left_at"])
CallEvent.objects.create(
session=session,
event_type=CallEventType.JOIN,
actor=self.user,
target_user=self.user,
)
return participant
@database_sync_to_async
def _record_leave(self):
try:
session = CallSession.objects.get(pk=self.session_id)
participant = CallParticipant.objects.get(session=session, user=self.user)
participant.status = CallParticipantStatus.LEFT
participant.left_at = now()
participant.save(update_fields=["status", "left_at"])
CallEvent.objects.create(
session=session,
event_type=CallEventType.LEAVE,
actor=self.user,
target_user=self.user,
)
except Exception:
pass
@database_sync_to_async
def _get_active_participants(self):
session = CallSession.objects.get(pk=self.session_id)
return [
{
"user_id": str(p.user_id),
"user_name": p.user.get_full_name() or str(p.user),
"is_muted": p.is_muted,
}
for p in session.call_participants.filter(
status=CallParticipantStatus.JOINED,
).select_related("user")
]
@database_sync_to_async
def _log_moderation_event(self, event_type, target_user_id):
from django.contrib.auth import get_user_model
User = get_user_model()
session = CallSession.objects.get(pk=self.session_id)
try:
target = User.objects.get(pk=target_user_id)
except User.DoesNotExist:
target = None
CallEvent.objects.create(
session=session,
event_type=event_type,
actor=self.user,
target_user=target,
)
@database_sync_to_async
def _set_participant_muted(self, user_id, is_muted):
try:
session = CallSession.objects.get(pk=self.session_id)
p = CallParticipant.objects.get(
session=session,
user_id=user_id,
status=CallParticipantStatus.JOINED,
)
p.is_muted = is_muted
p.save(update_fields=["is_muted"])
except CallParticipant.DoesNotExist:
pass
@database_sync_to_async
def _is_instructor(self):
return _user_is_instructor(self.user)
async def _send_error(self, code):
await self.send(text_data=json.dumps({"type": "error", "code": code}))
def _user_is_instructor(user):
try:
from accounts.services import get_effective_capabilities
return "training:signoff" in get_effective_capabilities(user)
except Exception:
return False