Spaces:
Sleeping
Sleeping
| """ | |
| نظام أولوية المالك - 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__) | |
| 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) | |