forked from hesabix/arc
349 lines
13 KiB
Python
Executable file
349 lines
13 KiB
Python
Executable file
"""
|
||
Queue Service با استفاده از RQ (Redis Queue)
|
||
برای مدیریت background jobs
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import json
|
||
import logging
|
||
from typing import Any, Dict, Optional, Callable
|
||
from datetime import timedelta
|
||
|
||
import rq
|
||
from rq import Queue, Worker
|
||
from rq.job import Job
|
||
from rq.registry import FailedJobRegistry
|
||
from redis import Redis
|
||
|
||
from app.core.cache import get_redis_client
|
||
from app.core.settings import get_settings
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
# Queue names
|
||
QUEUE_DEFAULT = "default"
|
||
QUEUE_HIGH_PRIORITY = "high"
|
||
QUEUE_LOW_PRIORITY = "low"
|
||
QUEUE_EMAIL = "email"
|
||
QUEUE_REPORTS = "reports"
|
||
QUEUE_EXPORTS = "exports"
|
||
QUEUE_TAX = "tax" # Queue مخصوص عملیات مالیاتی
|
||
|
||
|
||
def load_redis_settings_from_configuration() -> tuple[bool, str, int, int, Optional[str]]:
|
||
"""
|
||
خواندن تنظیمات Redis از دیتابیس (تنظیمات کل سیستم) یا متغیرهای محیطی.
|
||
اگر مدیر Redis را غیرفعال کرده باشد، enabled=False برمیگردد.
|
||
"""
|
||
settings = get_settings()
|
||
redis_enabled = getattr(settings, "redis_enabled", False)
|
||
redis_host = getattr(settings, "redis_host", "localhost")
|
||
redis_port = getattr(settings, "redis_port", 6379)
|
||
redis_db = getattr(settings, "redis_db", 0)
|
||
redis_password = getattr(settings, "redis_password", None)
|
||
try:
|
||
from adapters.db.session import get_db
|
||
from app.services.system_settings_service import get_redis_configuration
|
||
|
||
db_gen = get_db()
|
||
db = next(db_gen)
|
||
try:
|
||
redis_config = get_redis_configuration(db)
|
||
redis_enabled = redis_config.get("enabled", False)
|
||
redis_host = redis_config.get("host", "localhost")
|
||
redis_port = redis_config.get("port", 6379)
|
||
redis_db = redis_config.get("db", 0)
|
||
redis_password = redis_config.get("password")
|
||
finally:
|
||
db.close()
|
||
except Exception:
|
||
pass
|
||
return bool(redis_enabled), str(redis_host), int(redis_port), int(redis_db), redis_password
|
||
|
||
|
||
def is_redis_enabled_in_configuration() -> bool:
|
||
"""آیا Redis در تنظیمات سیستم روشن است (بدون تست اتصال شبکه)."""
|
||
return load_redis_settings_from_configuration()[0]
|
||
|
||
|
||
def get_redis_connection() -> Optional[Redis]:
|
||
"""دریافت Redis connection برای RQ"""
|
||
client = get_redis_client()
|
||
if client is None:
|
||
return None
|
||
|
||
redis_enabled, redis_host, redis_port, redis_db, redis_password = load_redis_settings_from_configuration()
|
||
|
||
if not redis_enabled:
|
||
return None
|
||
|
||
try:
|
||
# RQ نیاز به Redis بدون decode_responses دارد
|
||
redis_conn = Redis(
|
||
host=redis_host,
|
||
port=redis_port,
|
||
db=redis_db,
|
||
password=redis_password,
|
||
decode_responses=False, # RQ نیاز به bytes دارد
|
||
socket_connect_timeout=2,
|
||
socket_timeout=2,
|
||
retry_on_timeout=True,
|
||
health_check_interval=30,
|
||
)
|
||
redis_conn.ping()
|
||
return redis_conn
|
||
except Exception as e:
|
||
logger.warning(f"Failed to create Redis connection for RQ: {e}")
|
||
return None
|
||
|
||
|
||
def get_queue(name: str = QUEUE_DEFAULT) -> Optional[Queue]:
|
||
"""دریافت queue با نام مشخص"""
|
||
redis_conn = get_redis_connection()
|
||
if redis_conn is None:
|
||
return None
|
||
return Queue(name, connection=redis_conn)
|
||
|
||
|
||
class QueueService:
|
||
"""سرویس مدیریت Queue"""
|
||
|
||
def __init__(self):
|
||
self.redis_conn = get_redis_connection()
|
||
self.enabled = self.redis_conn is not None
|
||
|
||
def enqueue(
|
||
self,
|
||
func: Callable,
|
||
*args,
|
||
queue_name: str = QUEUE_DEFAULT,
|
||
timeout: int = 300,
|
||
result_ttl: int = 3600,
|
||
job_id: Optional[str] = None,
|
||
**kwargs
|
||
) -> Optional[Job]:
|
||
"""
|
||
اضافه کردن job به queue
|
||
|
||
Args:
|
||
func: تابعی که باید اجرا شود
|
||
*args: آرگومانهای positional
|
||
queue_name: نام queue
|
||
timeout: حداکثر زمان اجرا (ثانیه)
|
||
result_ttl: زمان نگهداری نتیجه (ثانیه)
|
||
job_id: شناسه job (اختیاری)
|
||
**kwargs: آرگومانهای keyword
|
||
|
||
Returns:
|
||
Job object یا None در صورت خطا
|
||
"""
|
||
if not self.enabled:
|
||
logger.warning("Queue service is disabled (Redis not available)")
|
||
return None
|
||
|
||
try:
|
||
queue = get_queue(queue_name)
|
||
if queue is None:
|
||
return None
|
||
|
||
job = queue.enqueue(
|
||
func,
|
||
*args,
|
||
job_id=job_id,
|
||
timeout=timeout,
|
||
result_ttl=result_ttl,
|
||
**kwargs
|
||
)
|
||
logger.info(f"Job {job.id} enqueued to {queue_name} queue")
|
||
return job
|
||
except Exception as e:
|
||
logger.error(f"Error enqueueing job: {e}", exc_info=True)
|
||
return None
|
||
|
||
def get_job(self, job_id: str) -> Optional[Job]:
|
||
"""دریافت job با شناسه"""
|
||
if not self.enabled:
|
||
logger.debug(f"Queue service disabled, cannot fetch job {job_id}")
|
||
return None
|
||
|
||
try:
|
||
job = Job.fetch(job_id, connection=self.redis_conn)
|
||
logger.debug(f"Job {job_id} found in Redis")
|
||
return job
|
||
except rq.exceptions.NoSuchJobError:
|
||
# Job وجود ندارد - این یک وضعیت عادی است
|
||
logger.debug(f"Job {job_id} not found in Redis (NoSuchJobError)")
|
||
return None
|
||
except Exception as e:
|
||
logger.warning(f"Error fetching job {job_id}: {e}", exc_info=True)
|
||
return None
|
||
|
||
def get_job_status(self, job_id: str) -> Optional[Dict[str, Any]]:
|
||
"""دریافت وضعیت job به صورت dictionary"""
|
||
logger.info(f"[QUEUE_SERVICE] get_job_status called for job_id: {job_id}")
|
||
job = self.get_job(job_id)
|
||
if job is None:
|
||
logger.info(f"[QUEUE_SERVICE] Job {job_id} not found in main queue, checking failed registry")
|
||
# اگر job در queue اصلی پیدا نشد، در failed registry بررسی کن
|
||
failed_status = self._get_failed_job_status(job_id)
|
||
if failed_status:
|
||
logger.info(f"[QUEUE_SERVICE] Job {job_id} found in failed registry")
|
||
else:
|
||
logger.info(f"[QUEUE_SERVICE] Job {job_id} not found in failed registry either")
|
||
return failed_status
|
||
|
||
try:
|
||
status = {
|
||
"id": job.id,
|
||
"state": job.get_status(),
|
||
"created_at": job.created_at.isoformat() if job.created_at else None,
|
||
"started_at": job.started_at.isoformat() if job.started_at else None,
|
||
"ended_at": job.ended_at.isoformat() if job.ended_at else None,
|
||
"result": None,
|
||
"error": None,
|
||
"updated_at": (job.ended_at or job.started_at or job.created_at).isoformat() if (job.ended_at or job.started_at or job.created_at) else None,
|
||
}
|
||
|
||
# اضافه کردن نتیجه یا خطا
|
||
if job.is_finished:
|
||
try:
|
||
status["result"] = job.result
|
||
except Exception:
|
||
pass
|
||
elif job.is_failed:
|
||
try:
|
||
if job.exc_info:
|
||
status["error"] = str(job.exc_info)
|
||
else:
|
||
# اگر exc_info نباشد، سعی میکنیم از exception message استفاده کنیم
|
||
try:
|
||
status["error"] = str(job.exc_info) if hasattr(job, 'exc_info') and job.exc_info else "Job failed without error details"
|
||
except Exception:
|
||
status["error"] = "Job failed without error details"
|
||
except Exception as e:
|
||
logger.warning(f"Error getting job error info for {job_id}: {e}")
|
||
status["error"] = "Job failed (error details unavailable)"
|
||
|
||
# اضافه کردن metadata اگر وجود دارد
|
||
if job.meta:
|
||
status["meta"] = job.meta
|
||
|
||
return status
|
||
except Exception as e:
|
||
logger.error(f"Error getting job status {job_id}: {e}", exc_info=True)
|
||
return None
|
||
|
||
def cancel_job(self, job_id: str) -> bool:
|
||
"""لغو job"""
|
||
job = self.get_job(job_id)
|
||
if job is None:
|
||
return False
|
||
|
||
try:
|
||
job.cancel()
|
||
logger.info(f"Job {job_id} cancelled")
|
||
return True
|
||
except Exception as e:
|
||
logger.error(f"Error cancelling job {job_id}: {e}", exc_info=True)
|
||
return False
|
||
|
||
def delete_job(self, job_id: str) -> bool:
|
||
"""حذف job"""
|
||
job = self.get_job(job_id)
|
||
if job is None:
|
||
return False
|
||
|
||
try:
|
||
job.delete()
|
||
logger.info(f"Job {job_id} deleted")
|
||
return True
|
||
except Exception as e:
|
||
logger.error(f"Error deleting job {job_id}: {e}", exc_info=True)
|
||
return False
|
||
|
||
def get_queue_length(self, queue_name: str = QUEUE_DEFAULT) -> int:
|
||
"""دریافت تعداد jobs در queue"""
|
||
queue = get_queue(queue_name)
|
||
if queue is None:
|
||
return 0
|
||
return len(queue)
|
||
|
||
def _get_failed_job_status(self, job_id: str) -> Optional[Dict[str, Any]]:
|
||
"""بررسی job در failed registry"""
|
||
if not self.enabled:
|
||
return None
|
||
|
||
try:
|
||
# بررسی در تمام queue ها
|
||
for queue_name in [QUEUE_DEFAULT, QUEUE_REPORTS, QUEUE_HIGH_PRIORITY,
|
||
QUEUE_LOW_PRIORITY, QUEUE_EMAIL, QUEUE_EXPORTS, QUEUE_TAX]:
|
||
try:
|
||
queue = Queue(queue_name, connection=self.redis_conn)
|
||
failed_registry = FailedJobRegistry(queue=queue)
|
||
if job_id in failed_registry:
|
||
# job در failed registry است
|
||
try:
|
||
job = Job.fetch(job_id, connection=self.redis_conn)
|
||
if job:
|
||
error_msg = "Job failed"
|
||
try:
|
||
if job.exc_info:
|
||
error_msg = str(job.exc_info)
|
||
elif hasattr(job, 'exc_info') and job.exc_info:
|
||
error_msg = str(job.exc_info)
|
||
except Exception:
|
||
error_msg = "Job failed (error details unavailable)"
|
||
|
||
return {
|
||
"id": job.id,
|
||
"state": "failed",
|
||
"created_at": job.created_at.isoformat() if job.created_at else None,
|
||
"started_at": job.started_at.isoformat() if job.started_at else None,
|
||
"ended_at": job.ended_at.isoformat() if job.ended_at else None,
|
||
"result": None,
|
||
"error": error_msg,
|
||
"updated_at": (job.ended_at or job.started_at or job.created_at).isoformat() if (job.ended_at or job.started_at or job.created_at) else None,
|
||
"meta": job.meta if job.meta else None,
|
||
}
|
||
except rq.exceptions.NoSuchJobError:
|
||
# Job در failed registry است اما از Redis حذف شده
|
||
logger.debug(f"Job {job_id} is in failed registry but not in Redis")
|
||
continue
|
||
except Exception as e:
|
||
logger.debug(f"Error checking queue {queue_name} for failed job {job_id}: {e}")
|
||
continue
|
||
except Exception as e:
|
||
logger.debug(f"Error checking failed registry for job {job_id}: {e}")
|
||
|
||
return None
|
||
|
||
def get_failed_jobs(self, limit: int = 10) -> list[Dict[str, Any]]:
|
||
"""دریافت لیست jobs ناموفق"""
|
||
if not self.enabled:
|
||
return []
|
||
|
||
try:
|
||
failed_registry = FailedJobRegistry(queue=Queue(connection=self.redis_conn))
|
||
job_ids = failed_registry.get_job_ids(0, limit - 1)
|
||
jobs = []
|
||
for job_id in job_ids:
|
||
status = self.get_job_status(job_id)
|
||
if status:
|
||
jobs.append(status)
|
||
return jobs
|
||
except Exception as e:
|
||
logger.error(f"Error getting failed jobs: {e}", exc_info=True)
|
||
return []
|
||
|
||
|
||
_queue_service: Optional[QueueService] = None
|
||
|
||
|
||
def get_queue_service() -> QueueService:
|
||
"""دریافت instance از QueueService"""
|
||
global _queue_service
|
||
if _queue_service is None:
|
||
_queue_service = QueueService()
|
||
return _queue_service
|
||
|