import json import asyncio from channels.generic.websocket import AsyncWebsocketConsumer from channels.db import database_sync_to_async from django.utils import timezone from datetime import timedelta from django.urls import reverse from .models import QuizGame, QuizGameParticipant class GameConsumer(AsyncWebsocketConsumer): async def connect(self): self.join_code = self.scope['url_route']['kwargs']['join_code'] self.game_group_name = f'game_{self.join_code}' # Join game group await self.channel_layer.group_add( self.game_group_name, self.channel_name ) await self.accept() async def disconnect(self, close_code): # Leave game group await self.channel_layer.group_discard( self.game_group_name, self.channel_name ) async def receive(self, text_data): text_data_json = json.loads(text_data) message_type = text_data_json['type'] if message_type == 'submit_answer': participant_id = text_data_json['participant_id'] answer = text_data_json['answer'] time_remaining = text_data_json.get('time_remaining', 0) # Save answer and update participant score await self.save_answer(participant_id, answer, time_remaining) # Notify host about the new answer await self.channel_layer.group_send( self.game_group_name, { 'type': 'participant_answer', 'answer': answer } ) elif message_type == 'get_answer_stats': stats = await self.get_answer_stats() await self.send(text_data=json.dumps({ 'type': 'answer_stats', 'stats': stats })) elif message_type == 'update_participants': participants = await self.get_participants() await self.channel_layer.group_send( self.game_group_name, { 'type': 'participant_list_update', 'participants': participants } ) elif message_type == 'next_question': host_id = text_data_json['host_id'] if await self.verify_host(host_id): await self.advance_to_next_question() elif message_type == 'finish_game': host_id = text_data_json['host_id'] if await self.verify_host(host_id): await self.finish_game() elif message_type == 'advance_to_scores': host_id = text_data_json['host_id'] if await self.verify_host(host_id): await self.show_scores() async def participant_answer(self, event): await self.send(text_data=json.dumps({ 'type': 'participant_answer', 'answer': event['answer'] })) async def participant_list_update(self, event): await self.send(text_data=json.dumps({ 'type': 'participant_list_update', 'participants': event['participants'] })) async def game_state_update(self, event): await self.send(text_data=json.dumps({ 'type': 'game_state_update', 'action': event['action'], 'redirect_url': event.get('redirect_url') })) @database_sync_to_async def verify_host(self, host_id): try: quiz_game = QuizGame.objects.get(join_code=self.join_code) return quiz_game.host_id == host_id except QuizGame.DoesNotExist: return False @database_sync_to_async def save_answer(self, participant_id, answer_index, time_remaining): try: participant = QuizGameParticipant.objects.get( quiz_game__join_code=self.join_code, participant_id=participant_id ) quiz_game = participant.quiz_game current_question = quiz_game.quiz_id.questions.all()[quiz_game.current_question_index] # Check if answer is correct is_correct = False if answer_index >= 0: # -1 means timeout is_correct = current_question.data['options'][answer_index]['is_correct'] # Calculate score based on correctness and time score = 0 if is_correct: base_score = 1000 time_factor = time_remaining / 30000 # 30 seconds max score = int(base_score * (0.5 + 0.5 * time_factor)) participant.last_score=score participant.score += score participant.last_answer_correct = is_correct participant.save() return True except (QuizGameParticipant.DoesNotExist, IndexError): return False @database_sync_to_async def get_answer_stats(self): try: quiz_game = QuizGame.objects.get(join_code=self.join_code) participants = QuizGameParticipant.objects.filter(quiz_game=quiz_game) # Count answers for each option stats = {} for participant in participants: if participant.last_answer is not None: stats[participant.last_answer] = stats.get(participant.last_answer, 0) + 1 return stats except QuizGame.DoesNotExist: return {} @database_sync_to_async def get_participants(self): try: quiz_game = QuizGame.objects.get(join_code=self.join_code) participants = QuizGameParticipant.objects.filter(quiz_game=quiz_game) return [ { 'id': p.participant_id, 'display_name': p.display_name, 'score': p.score } for p in participants ] except QuizGame.DoesNotExist: return [] @database_sync_to_async def _advance_to_next_question(self): try: quiz_game = QuizGame.objects.get(join_code=self.join_code) quiz = quiz_game.quiz_id # Move to next question quiz_game.current_question_index += 1 # Check if we've reached the end if quiz_game.current_question_index >= quiz.questions.count(): quiz_game.current_state = 'finished' quiz_game.save() return 'finished' else: # Update game state and start time quiz_game.current_state = 'question' quiz_game.question_start_time = timezone.now() quiz_game.save() return 'question' except QuizGame.DoesNotExist: return None async def advance_to_next_question(self): result = await self._advance_to_next_question() if result == 'finished': # Notify clients to redirect to finished page await self.channel_layer.group_send( self.game_group_name, { 'type': 'game_state_update', 'action': 'finish_game', 'redirect_url': reverse('play:finished', kwargs={'join_code': self.join_code}) } ) elif result == 'question': # Notify clients to redirect to next question await self.channel_layer.group_send( self.game_group_name, { 'type': 'game_state_update', 'action': 'next_question', 'redirect_url': reverse('play:question', kwargs={'join_code': self.join_code}) } ) @database_sync_to_async def _show_scores(self): try: quiz_game = QuizGame.objects.get(join_code=self.join_code) quiz_game.current_state = 'scores' quiz_game.save() return True except QuizGame.DoesNotExist: return False async def show_scores(self): success = await self._show_scores() if success: # Notify clients to redirect to scores page await self.channel_layer.group_send( self.game_group_name, { 'type': 'game_state_update', 'action': 'show_scores', 'redirect_url': reverse('play:scores', kwargs={'join_code': self.join_code}) } ) @database_sync_to_async def _finish_game(self): try: quiz_game = QuizGame.objects.get(join_code=self.join_code) quiz_game.current_state = 'finished' quiz_game.save() return True except QuizGame.DoesNotExist: return False async def finish_game(self): success = await self._finish_game() if success: # Notify clients to redirect to finished page await self.channel_layer.group_send( self.game_group_name, { 'type': 'game_state_update', 'action': 'finish_game', 'redirect_url': reverse('play:finished', kwargs={'join_code': self.join_code}) } ) class LobbyConsumer(AsyncWebsocketConsumer): async def connect(self): self.join_code = self.scope['url_route']['kwargs']['join_code'] self.room_group_name = f'game_{self.join_code}' self.heartbeat_task = None # Join room group await self.channel_layer.group_add( self.room_group_name, self.channel_name ) await self.accept() # Start heartbeat check self.heartbeat_task = asyncio.create_task(self.check_participants_heartbeat()) # Send current participants list participants = await self.get_participants() await self.send(text_data=json.dumps({ 'type': 'participant_list', 'participants': participants })) async def disconnect(self, close_code): # Cancel heartbeat task if self.heartbeat_task: self.heartbeat_task.cancel() try: await self.heartbeat_task except asyncio.CancelledError: pass # Leave room group await self.channel_layer.group_discard( self.room_group_name, self.channel_name ) async def receive(self, text_data): text_data_json = json.loads(text_data) message_type = text_data_json['type'] if message_type == 'heartbeat': # Update participant's last activity timestamp participant_id = text_data_json.get('participant_id') if participant_id: await self._update_participant_heartbeat(participant_id) elif message_type == 'update_participants': participants = await self.get_participants() await self.channel_layer.group_send( self.room_group_name, { 'type': 'participant_list_update', 'participants': participants } ) elif message_type == 'start_game': # Verify sender is host sender_id = text_data_json.get('host_id') if await self._is_host(sender_id): # Update game state and notify all participants await self._start_game() await self.channel_layer.group_send( self.room_group_name, { 'type': 'game_state_update', 'action': 'start_game', 'redirect_url': f'/play/game/{self.join_code}/0' } ) async def game_state_update(self, event): # Send game state update to WebSocket await self.send(text_data=json.dumps({ 'type': 'game_state_update', 'action': event['action'], 'redirect_url': event.get('redirect_url') })) @database_sync_to_async def _start_game(self): game = QuizGame.objects.get(join_code=self.join_code) game.current_state = 'question' game.current_question_index = 0 game.question_start_time = timezone.now() game.save() async def receive(self, text_data): text_data_json = json.loads(text_data) message_type = text_data_json['type'] if message_type == 'heartbeat': # Update participant's last activity timestamp participant_id = text_data_json.get('participant_id') if participant_id: await self._update_participant_heartbeat(participant_id) elif message_type == 'update_participants': participants = await self.get_participants() await self.channel_layer.group_send( self.room_group_name, { 'type': 'participant_list_update', 'participants': participants } ) elif message_type == 'start_game': # Verify sender is host sender_id = text_data_json.get('host_id') if await self._is_host(sender_id): # Update game state and notify all participants await self._start_game() await self.channel_layer.group_send( self.room_group_name, { 'type': 'game_state_update', 'action': 'start_game', 'redirect_url': f'/play/game/{self.join_code}/0' } ) elif message_type == 'submit_answer': participant_id = text_data_json.get('participant_id') answer = text_data_json.get('answer') time_remaining = text_data_json.get('time_remaining') if participant_id: score = await self._process_answer(participant_id, answer, time_remaining) await self.send(text_data=json.dumps({ 'type': 'answer_processed', 'score': score })) elif message_type == 'next_question': sender_id = text_data_json.get('host_id') if await self._is_host(sender_id): next_state = await self._advance_game_state() await self.channel_layer.group_send( self.room_group_name, { 'type': 'game_state_update', 'action': next_state['action'], 'redirect_url': next_state['redirect_url'] } ) async def participant_list_update(self, event): participants = event['participants'] await self.send(text_data=json.dumps({ 'type': 'participant_list', 'participants': participants })) @database_sync_to_async def get_participants(self): game = QuizGame.objects.get(join_code=self.join_code) participants = QuizGameParticipant.objects.filter(quiz_game=game) return [{'id': str(p.participant_id), 'name': p.display_name} for p in participants] @database_sync_to_async def _update_participant_heartbeat(self, participant_id): try: participant = QuizGameParticipant.objects.get(participant_id=participant_id) participant.save() # This will update last_heartbeat due to auto_now=True except QuizGameParticipant.DoesNotExist: pass @database_sync_to_async def _is_host(self, host_id): try: return QuizGame.objects.filter(join_code=self.join_code, host_id=host_id).exists() except QuizGame.DoesNotExist: return False @database_sync_to_async def _start_game(self): game = QuizGame.objects.get(join_code=self.join_code) game.current_state = 'question' game.current_question_index = 0 game.question_start_time = timezone.now() game.save() @database_sync_to_async def _process_answer(self, participant_id, answer, time_remaining): game = QuizGame.objects.get(join_code=self.join_code) participant = QuizGameParticipant.objects.get(participant_id=participant_id) question = game.quiz_id.questions.all()[game.current_question_index] # Load question data question_data = json.loads(question.data) correct_answer = question_data.get('correct_answer') # Calculate score based on correctness and time score = 0 if answer == correct_answer: # Base score for correct answer + bonus for speed score = 1000 + int(time_remaining * 10) # 10 points per remaining second participant.last_score= score participant.score += score participant.save() return score @database_sync_to_async def _advance_game_state(self): game = QuizGame.objects.get(join_code=self.join_code) total_questions = game.quiz_id.questions.count() if game.current_state == 'question': game.current_state = 'scores' game.save() return { 'action': 'show_scores', 'redirect_url': f'/play/game/{self.join_code}/scores' } elif game.current_state == 'scores': if game.current_question_index + 1 < total_questions: game.current_state = 'question' game.current_question_index += 1 game.question_start_time = timezone.now() game.save() return { 'action': 'next_question', 'redirect_url': f'/play/game/{self.join_code}/{game.current_question_index}' } else: game.current_state = 'finished' game.save() return { 'action': 'game_finished', 'redirect_url': f'/play/game/{self.join_code}/finished' } @database_sync_to_async def _remove_inactive_participants(self): # Remove participants who haven't sent a heartbeat in the last 30 seconds timeout = timezone.now() - timedelta(seconds=30) quiz_game = QuizGame.objects.get(join_code=self.join_code) QuizGameParticipant.objects.filter( quiz_game=quiz_game, last_heartbeat__lt=timeout ).delete() async def check_participants_heartbeat(self): while True: try: await asyncio.sleep(10) # Check every 10 seconds await self._remove_inactive_participants() # Send updated participant list participants = await self.get_participants() await self.channel_layer.group_send( self.room_group_name, { 'type': 'participant_list_update', 'participants': participants } ) except asyncio.CancelledError: break except Exception as e: print(f'Error in heartbeat check: {e}') await asyncio.sleep(10) # Wait before retrying