zai-telegram-bot / owner_priority.py
Ejdjdososs's picture
Upload owner_priority.py with huggingface_hub
e94914b verified
Raw
History Blame Contribute Delete
5.15 kB
"""
نظام أولوية المالك - 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)