""" Moodle Queue Service (admin-edu-space) Provides asynchronous, fault-tolerant queuing and processing of synchronization operations between admin-edu-space and Moodle 4.1. Implements Exponential Backoff and Dead Letter Queue (DLQ). """ import logging from datetime import datetime, timedelta from app import db from app.models.sync_task import MoodleSyncTask from app.services.moodle_client import moodle_client logger = logging.getLogger(__name__) class MoodleQueueService: @staticmethod def enqueue_task(action: str, entity_type: str, entity_id: str = None, payload: dict = None, max_attempts: int = 5) -> MoodleSyncTask: """ Enqueues a new synchronization task to be processed asynchronously. Guarantees that local transactions are never blocked by Moodle unavailability. """ task = MoodleSyncTask( action=action, entity_type=entity_type, entity_id=str(entity_id) if entity_id else None, payload=payload or {}, status='PENDING', attempts=0, max_attempts=max_attempts, next_retry_at=datetime.utcnow() ) db.session.add(task) db.session.commit() logger.info(f"[MoodleQueue] Enqueued task {task.id}: action={action}, entity={entity_type}:{entity_id}") return task @staticmethod def process_pending_tasks(batch_size: int = 20) -> dict: """ Processes a batch of pending or retrying tasks. Uses exponential backoff for retries and sends to Dead Letter Queue (FAILED) if max_attempts are exceeded. """ now = datetime.utcnow() tasks = MoodleSyncTask.query.filter( MoodleSyncTask.status.in_(['PENDING', 'RETRYING']), (MoodleSyncTask.next_retry_at == None) | (MoodleSyncTask.next_retry_at <= now) ).order_by(MoodleSyncTask.created_at.asc()).limit(batch_size).all() results = { 'processed': 0, 'succeeded': 0, 'failed': 0, 'retrying': 0 } if not tasks: return results for task in tasks: results['processed'] += 1 task.status = 'PROCESSING' task.updated_at = datetime.utcnow() db.session.commit() try: MoodleQueueService._execute_task_action(task) task.status = 'COMPLETED' task.error_message = None task.updated_at = datetime.utcnow() db.session.commit() results['succeeded'] += 1 logger.info(f"[MoodleQueue] Task {task.id} ({task.action}) completed successfully.") except Exception as e: task.attempts += 1 task.error_message = str(e) task.updated_at = datetime.utcnow() if task.attempts >= task.max_attempts: task.status = 'FAILED' # Dead Letter Queue (DLQ) task.next_retry_at = None results['failed'] += 1 logger.error(f"[MoodleQueue] Task {task.id} permanently failed (DLQ): {e}") else: task.status = 'RETRYING' # Exponential Backoff: 30s, 60s, 120s, 240s... capped at 1 hour delay_seconds = min(3600, (2 ** task.attempts) * 30) task.next_retry_at = datetime.utcnow() + timedelta(seconds=delay_seconds) results['retrying'] += 1 logger.warning(f"[MoodleQueue] Task {task.id} failed attempt {task.attempts}/{task.max_attempts}. Next retry in {delay_seconds}s: {e}") db.session.commit() return results @staticmethod def _execute_task_action(task: MoodleSyncTask): """ Executes the specific Moodle Web Service operation. Raises an exception on failure or error response. """ payload = task.payload or {} action = task.action.upper() if action == 'CREATE_USER': users = payload.get('users') or [payload] res = moodle_client.create_users(users) return res elif action == 'UPDATE_USER': users = payload.get('users') or [payload] res = moodle_client.update_users(users) return res elif action == 'ENROL_USER': enrolments = payload.get('enrolments') or [payload] res = moodle_client.enrol_users(enrolments) return res elif action == 'UNENROL_USER': enrolments = payload.get('enrolments') or [payload] res = moodle_client.unenrol_users(enrolments) return res elif action == 'ASSIGN_ROLE': role_id = payload.get('role_id') user_id = payload.get('user_id') context_id = payload.get('context_id', 1) res = moodle_client.assign_role(role_id, user_id, context_id) return res elif action == 'UNASSIGN_ROLE': role_id = payload.get('role_id') user_id = payload.get('user_id') context_id = payload.get('context_id', 1) res = moodle_client.unassign_role(role_id, user_id, context_id) return res else: raise ValueError(f"Unsupported sync action: {action}") @staticmethod def retry_task(task_id: int) -> bool: """ Manually re-enqueues a task from DLQ or error state back to PENDING. """ task = MoodleSyncTask.query.get(task_id) if not task: return False task.status = 'PENDING' task.attempts = 0 task.next_retry_at = datetime.utcnow() task.error_message = None task.updated_at = datetime.utcnow() db.session.commit() logger.info(f"[MoodleQueue] Task {task_id} manually reset to PENDING.") return True @staticmethod def retry_all_failed() -> int: """ Retries all tasks in FAILED status (DLQ). """ failed_tasks = MoodleSyncTask.query.filter_by(status='FAILED').all() count = 0 for task in failed_tasks: task.status = 'PENDING' task.attempts = 0 task.next_retry_at = datetime.utcnow() task.updated_at = datetime.utcnow() count += 1 db.session.commit() logger.info(f"[MoodleQueue] Reset {count} failed tasks back to PENDING.") return count @staticmethod def get_queue_summary() -> dict: """ Returns stats about tasks currently in the queue. """ counts = { 'PENDING': 0, 'PROCESSING': 0, 'RETRYING': 0, 'COMPLETED': 0, 'FAILED': 0 } from sqlalchemy import func rows = db.session.query(MoodleSyncTask.status, func.count(MoodleSyncTask.id)).group_by(MoodleSyncTask.status).all() for status, count in rows: if status in counts: counts[status] = count total = sum(counts.values()) return { 'summary': counts, 'total': total, 'pending_total': counts['PENDING'] + counts['RETRYING'] + counts['PROCESSING'], 'failed_dlq': counts['FAILED'] } moodle_queue_service = MoodleQueueService()