""" نظام أولوية المالك - Priority Queue - رسائل المالك لها الأولوية القصوى - تُعالَج قبل رسائل المستخدمين الآخرين حتى لو وصلت متأخرة """ import asyncio import logging import time from dataclasses import dataclass, field from typing import Any, Callable, Optional from config import config logger = logging.getLogger(__name__) @dataclass(order=True) class PriorityTask: """مهمة في طابور الأولوية - أقل priority = أعلى أولوية""" priority: int # 0 = owner (highest), 10 = regular user sequence: int # ترتيب الوصول (للـ FIFO داخل نفس الأولوية) task_id: str = field(compare=False) # معرّف فريد func: Callable = field(compare=False) # الدالة المطلوب تنفيذها args: tuple = field(compare=False, default=()) kwargs: dict = field(compare=False, default_factory=dict) class PriorityProcessor: """معالج طابور الأولوية""" def __init__(self, max_workers: int = 4): self.queue: asyncio.PriorityQueue = asyncio.PriorityQueue() self.max_workers = max_workers self._workers: list = [] self._sequence = 0 self._running = False self._counter_lock = asyncio.Lock() self._task_results: dict = {} # task_id -> result self._task_events: dict = {} # task_id -> asyncio.Event async def start(self): """بدء عمال المعالجة""" if self._running: return self._running = True for i in range(self.max_workers): worker = asyncio.create_task(self._worker(i)) self._workers.append(worker) logger.info(f"Started {self.max_workers} priority queue workers") async def stop(self): """إيقاف العمال""" self._running = False # Send sentinel tasks to wake up workers for _ in self._workers: await self.queue.put(PriorityTask( priority=999, sequence=0, task_id="__stop__", func=lambda: None, )) for w in self._workers: try: await asyncio.wait_for(w, timeout=5) except (asyncio.TimeoutError, asyncio.CancelledError): w.cancel() self._workers.clear() async def _next_sequence(self) -> int: async with self._counter_lock: self._sequence += 1 return self._sequence async def submit( self, user_id: int, func: Callable, *args, task_id: str = "", **kwargs, ) -> Any: """ إرسال مهمة للطابور. - user_id يحدد الأولوية (المالك = 0، مشرف = 2، مستخدم = 10) - يعيد نتيجة الدالة مباشرة """ priority = config.priority_for(user_id) seq = await self._next_sequence() if not task_id: task_id = f"task-{seq}" event = asyncio.Event() self._task_events[task_id] = event task = PriorityTask( priority=priority, sequence=seq, task_id=task_id, func=func, args=args, kwargs=kwargs, ) await self.queue.put(task) logger.info(f"Submitted task {task_id} for user {user_id} with priority {priority}") # انتظر النتيجة await event.wait() result = self._task_results.pop(task_id, None) self._task_events.pop(task_id, None) if isinstance(result, Exception): raise result return result async def _worker(self, worker_id: int): """عامل معالجة الطابور""" logger.info(f"Priority queue worker {worker_id} started") while self._running: try: task: PriorityTask = await asyncio.wait_for(self.queue.get(), timeout=1.0) except asyncio.TimeoutError: continue if task.task_id == "__stop__": self.queue.task_done() break try: logger.debug( f"Worker {worker_id} processing task {task.task_id} " f"(priority={task.priority})" ) if asyncio.iscoroutinefunction(task.func): result = await task.func(*task.args, **task.kwargs) else: result = await asyncio.to_thread(task.func, *task.args, **task.kwargs) self._task_results[task.task_id] = result except Exception as e: logger.error(f"Task {task.task_id} failed: {e}", exc_info=True) self._task_results[task.task_id] = e finally: event = self._task_events.get(task.task_id) if event: event.set() self.queue.task_done() # Singleton priority_processor = PriorityProcessor(max_workers=4)