204 lines
7.2 KiB
Python
204 lines
7.2 KiB
Python
"""
|
|
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()
|